diff --git a/docker/mistserver.json b/docker/mistserver.json index dfb0f830..580b5511 100644 --- a/docker/mistserver.json +++ b/docker/mistserver.json @@ -115,16 +115,6 @@ "stream": { "debug": 5, "name": "stream", - "processes": [ - { - "debug": 5, - "exec": "streamplace live $wildcard", - "exit_unmask": false, - "inconsequential": false, - "process": "MKVExec", - "restart_type": "fixed" - } - ], "source": "push://", "stop_sessions": false, "tags": [] diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index 1b92b862..bde37570 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -10,7 +10,6 @@ import ( "net/http" "net/http/pprof" "os" - "regexp" "runtime" rtpprof "runtime/pprof" "strconv" @@ -45,19 +44,16 @@ func (a *StreamplaceAPI) ServeInternalHTTP(ctx context.Context) error { }) } -// lightweight way to authenticate push requests to ourself -var mkvRE *regexp.Regexp - -func init() { - mkvRE = regexp.MustCompile(`^\d+\.mkv$`) -} - func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, error) { router := httprouter.New() broker := misttriggers.NewTriggerBroker() + // serverCtx outlives any single trigger request — the Mist pull ingests + // spawned below run for the life of their stream, not the life of the + // PUSH_REWRITE request that announced it. + serverCtx := ctx broker.OnPushRewrite(func(ctx context.Context, payload *misttriggers.PushRewritePayload) (string, error) { - log.Log(ctx, "got push out start", "streamName", payload.StreamName, "url", payload.URL.String()) + log.Log(ctx, "got push rewrite", "streamName", payload.StreamName, "url", payload.URL.String()) // Extract the last part of the URL path urlPath := payload.URL.Path parts := strings.Split(urlPath, "/") @@ -77,6 +73,20 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err a.SignerCacheMu.Unlock() log.Log(ctx, "added key to cache", "mist-stream", out, "streamer", mediaSigner.Streamer()) + // The push is authed and named — ingest it by pulling Mist's live fMP4 + // output for the stream we just named. This replaces the old Mist-side + // MKVExec process (`streamplace live` POSTing MKV back to /live): fMP4 + // carries real decode timestamps, so ingest no longer reconstructs DTS. + // Mist accepts the push right after this trigger returns, so the pull + // retries briefly while the stream boots (mistPullConnect). + go func() { + if perr := a.MediaManager.MistPullIngest(serverCtx, out, mediaSigner); perr != nil { + log.Error(serverCtx, "mist pull ingest ended", "mist-stream", out, "streamer", mediaSigner.Streamer(), "error", perr) + } else { + log.Log(serverCtx, "mist pull ingest ended cleanly", "mist-stream", out, "streamer", mediaSigner.Streamer()) + } + }() + return out, nil }) triggerCollection := misttriggers.NewMistCallbackHandlersCollection(a.CLI, broker) @@ -209,21 +219,22 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err _, _ = io.ReadFull(bufrw.Reader, prebuf) } chunked := len(httpReq.TransferEncoding) > 0 && httpReq.TransferEncoding[0] == "chunked" - if derr := a.MediaManager.MKVIngestDetached(reqCtx, conn, prebuf, chunked, mediaSigner); derr != nil { + if derr := a.MediaManager.MP4IngestDetached(reqCtx, conn, prebuf, chunked, mediaSigner); derr != nil { log.Log(reqCtx, "isolated stream ended", "error", derr) } return // connection hijacked; the HTTP response is ours now } // The isolated path needs a hijackable HTTP/1.1 connection (which the - // only real MKV/RTMP-push client — a co-located MistServer pushing over - // localhost — always is). We don't support a non-hijack fallback: it - // couldn't receive mid-stream manifest updates and would stay stuck - // pre-live, so refuse. Such a client can use WHIP instead. + // real /live clients — the `streamplace live` CLI and tests pushing + // over localhost — always are; the Mist ingest itself now arrives via + // MistPullIngest, not this route). We don't support a non-hijack + // fallback: it couldn't receive mid-stream manifest updates and would + // stay stuck pre-live, so refuse. Such a client can use WHIP instead. log.Error(reqCtx, "isolated ingest requires a hijackable HTTP/1.1 connection; refusing push") errors.WriteHTTPInternalServerError(w, "isolated ingest requires a hijackable HTTP/1.1 connection; use WHIP", fmt.Errorf("connection is not hijackable")) return } else { - err = a.MediaManager.MKVIngest(reqCtx, r, mediaSigner) + err = a.MediaManager.MP4Ingest(reqCtx, r, mediaSigner) } if err != nil { @@ -234,7 +245,10 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err log.Log(reqCtx, "stream success", "url", httpReq.URL.String()) } - // route to accept an incoming mkv stream from OBS, segment it, and push the segments back to this HTTP handler + // route to accept an incoming fragmented-MP4 stream (the `streamplace live` + // CLI piping from stdin), segment it, and validate the signed segments. + // The co-located MistServer's streams are ingested by pulling its fMP4 + // output instead (MistPullIngest, kicked off from PUSH_REWRITE above). router.POST("/live/:key", handleIncomingStream) router.PUT("/live/:key", handleIncomingStream) diff --git a/pkg/cmd/live.go b/pkg/cmd/live.go index 1a814b16..6229ad59 100644 --- a/pkg/cmd/live.go +++ b/pkg/cmd/live.go @@ -8,7 +8,9 @@ import ( ) func Live(streamKey string, httpInternalAddr string) error { - // Create the URL for the live stream endpoint + // Live POSTs a fragmented-MP4 stream from stdin to the node's /live ingest + // route. (The Mist ingest bridge no longer uses this — the node pulls Mist's + // fMP4 output directly; see MistPullIngest.) url := fmt.Sprintf("http://%s/live/%s", httpInternalAddr, streamKey) // Create a new HTTP request with POST method @@ -18,7 +20,7 @@ func Live(streamKey string, httpInternalAddr string) error { } // Set appropriate headers if needed - req.Header.Set("Content-Type", "video/x-matroska") // Assuming MKV format, adjust if needed + req.Header.Set("Content-Type", "video/mp4") // fragmented MP4 from stdin // Create HTTP client and send the request client := &http.Client{} diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index fcdfdde9..9e58fb4c 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -866,8 +866,8 @@ func makeStreamCommand(build *config.BuildFlags) *urfavecli.Command { } // makeIngestWorkerCommand is the per-stream isolated ingest worker (Stage 1: -// MKV/RTMP push). The node spawns it; it is not meant for direct use. It reads -// the config handshake from fd 3, the MKV media from stdin, runs the mux + sign +// fMP4 / Mist pull). The node spawns it; it is not meant for direct use. It reads +// the config handshake from fd 3, the fragmented-MP4 media from stdin, runs the mux + sign // pipeline, and writes signed canonical .m4s frames to fd 4 — dedicated fds so // stray stdout/stderr can't corrupt the frame stream. A clean run ends with an // End frame; a fatal error emits an Error frame before exiting non-zero. @@ -920,7 +920,7 @@ func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { defer f.Close() raw = f } - return media.ServeMKVIngestWorkerSocket(ctx, cfg, media.WorkerInput(cfg, raw)) + return media.ServeMP4IngestWorkerSocket(ctx, cfg, media.WorkerInput(cfg, raw)) } framesFile := os.NewFile(4, "ingest-frames") @@ -933,7 +933,7 @@ func makeIngestWorkerCommand(build *config.BuildFlags) *urfavecli.Command { // This fd-4 path has no back-channel for manifest updates, so the // manifest stays whatever main built at spawn. It's not used in prod // (api requires a hijackable connection); kept for the worker self-test. - if err := media.RunMKVIngestWorker(ctx, cfg, os.Stdin, frames, func() []byte { return cfg.Manifest }); err != nil { + if err := media.RunMP4IngestWorker(ctx, cfg, os.Stdin, frames, func() []byte { return cfg.Manifest }); err != nil { _ = frames.Error(err.Error()) return err } @@ -995,7 +995,7 @@ func makeRTMPPushWorkerCommand(build *config.BuildFlags) *urfavecli.Command { func makeLiveCommand(build *config.BuildFlags) *urfavecli.Command { cli := config.CLI{Build: build} liveCmd := cli.NewCommand("live") - liveCmd.Usage = "start live stream" + liveCmd.Usage = "start live stream (pipe fragmented MP4 to stdin)" liveCmd.ArgsUsage = "[stream-key]" liveCmd.Action = func(ctx context.Context, cmd *urfavecli.Command) error { args := cmd.Args() diff --git a/pkg/media/frame_server.go b/pkg/media/frame_server.go index 6447e823..a60d7e38 100644 --- a/pkg/media/frame_server.go +++ b/pkg/media/frame_server.go @@ -27,7 +27,7 @@ const workerDrainGrace = 60 * time.Second // FrameWriter is the worker's segment sink. Stage 1 uses a direct framed pipe // (*ingestframe.Writer); the zero-downtime path uses *frameServer, which buffers // across a disconnected main and replays on reconnect. Structurally satisfied by -// *ingestframe.Writer, so RunMKVIngestWorker is agnostic to which it gets. +// *ingestframe.Writer, so RunMP4IngestWorker is agnostic to which it gets. type FrameWriter interface { Segment(seg []byte) error End() error @@ -165,15 +165,15 @@ func (s *frameServer) detachConn(conn net.Conn) { } } -// ServeMKVIngestWorkerSocket runs the ingest worker, delivering its signed +// ServeMP4IngestWorkerSocket runs the ingest worker, delivering its signed // segments to main over a per-session unix socket at cfg.SocketPath with // buffered reconnect — the zero-downtime path. It listens, serves the frame // stream (buffering across any main disconnect), runs the ingest, frames a // trailing End/Error, then lingers until main has drained the buffer before // removing the socket and returning. -func ServeMKVIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader) error { +func ServeMP4IngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader) error { if cfg.SocketPath == "" { - return fmt.Errorf("ServeMKVIngestWorkerSocket: empty socket path") + return fmt.Errorf("ServeMP4IngestWorkerSocket: empty socket path") } ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -192,7 +192,7 @@ func ServeMKVIngestWorkerSocket(ctx context.Context, cfg IngestWorkerConfig, std manifest := newManifestHolder(cfg.Manifest) go serveFrameSocket(ctx, ln, srv, manifest) - runErr := RunMKVIngestWorker(ctx, cfg, stdin, srv, manifest.get) + runErr := RunMP4IngestWorker(ctx, cfg, stdin, srv, manifest.get) if runErr != nil { _ = srv.Error(runErr.Error()) } else { diff --git a/pkg/media/frame_socket_e2e_test.go b/pkg/media/frame_socket_e2e_test.go index 33734803..2c201509 100644 --- a/pkg/media/frame_socket_e2e_test.go +++ b/pkg/media/frame_socket_e2e_test.go @@ -15,7 +15,7 @@ import ( ) // TestWorkerServesFramesOverSocket drives the zero-downtime transport end-to-end -// with a REAL ingest: ServeMKVIngestWorkerSocket runs the full mux+sign+transcode +// with a REAL ingest: ServeMP4IngestWorkerSocket runs the full mux+sign+transcode // pipeline and serves the resulting signed dual-codec segments over a unix // socket; a client connects and reads them through to a clean End. This proves // the socket path carries real signed media (the frameServer reconnect tests @@ -40,10 +40,10 @@ func TestWorkerServesFramesOverSocket(t *testing.T) { SocketPath: sock, } - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) serveDone := make(chan error, 1) - go func() { serveDone <- ServeMKVIngestWorkerSocket(ctx, cfg, bytes.NewReader(mkv)) }() + go func() { serveDone <- ServeMP4IngestWorkerSocket(ctx, cfg, bytes.NewReader(mp4)) }() // Connect once the worker's listener is up (retry the dial briefly). var conn net.Conn @@ -85,7 +85,7 @@ func TestWorkerServesFramesOverSocket(t *testing.T) { case serveErr := <-serveDone: require.NoError(t, serveErr) case <-time.After(30 * time.Second): - t.Fatal("ServeMKVIngestWorkerSocket did not return after the stream drained") + t.Fatal("ServeMP4IngestWorkerSocket did not return after the stream drained") } t.Logf("worker served %d signed segments + End over the socket", segs) } diff --git a/pkg/media/ingest_daemon.go b/pkg/media/ingest_daemon.go index 3cee3517..cab8825b 100644 --- a/pkg/media/ingest_daemon.go +++ b/pkg/media/ingest_daemon.go @@ -62,7 +62,7 @@ func SpawnIngestWorkerDetached(cfg IngestWorkerConfig, media *os.File) (*os.Proc // identifiable in a process listing; key material stays on fd 3. cmd := exec.Command(exe, "ingest-worker", cfg.StreamerDID) setDetached(cmd) // own session, survives a main restart (Linux) - // fd 3 = config; fd 4 = the fd-passed media connection (MKV/RTMP). WHIP owns + // fd 3 = config; fd 4 = the fd-passed media connection (the Mist fMP4 pull, or an fMP4 push). WHIP owns // its own PeerConnection, so it passes no media fd. cmd.ExtraFiles = []*os.File{cfgR} if media != nil { @@ -245,7 +245,7 @@ func (mm *MediaManager) ingestWorkerSocketDir() (string, error) { return dir, nil } -// MKVIngestDetached is the production zero-downtime entry: main has authed the +// MP4IngestDetached is the production zero-downtime entry: main has authed the // push and hijacked its connection; this fd-passes that connection to a DETACHED // worker (own session, survives a main restart) which ingests the media directly // and serves signed segments over a per-session unix socket, and then consumes @@ -256,7 +256,7 @@ func (mm *MediaManager) ingestWorkerSocketDir() (string, error) { // breaks the ingest nor loses output: the worker keeps signing into its buffer, // and the restarted main rediscovers the socket (DiscoverWorkerSockets) and // drains it. -func (mm *MediaManager) MKVIngestDetached(ctx context.Context, conn net.Conn, prebuf []byte, chunked bool, ms MediaSigner) error { +func (mm *MediaManager) MP4IngestDetached(ctx context.Context, conn net.Conn, prebuf []byte, chunked bool, ms MediaSigner) error { cfg, err := mm.buildWorkerConfig(ctx, ms) if err != nil { return err @@ -285,7 +285,7 @@ func (mm *MediaManager) MKVIngestDetached(ctx context.Context, conn net.Conn, pr if err != nil { return fmt.Errorf("spawn detached worker: %w", err) } - spmetrics.IngestWorkerStarts.WithLabelValues("mkv").Inc() + spmetrics.IngestWorkerStarts.WithLabelValues("mp4").Inc() // Ban / key revocation: the detached worker can't notice it itself (no // bus/model), so main watches and kills it. proc.Kill (not ctx cancel) so the @@ -301,7 +301,7 @@ func (mm *MediaManager) MKVIngestDetached(ctx context.Context, conn net.Conn, pr start := time.Now().UnixMilli() manifestSource := func() ([]byte, error) { return mm.streamerManifest(ctx, ms.Streamer(), start) } err = mm.ConsumeWorkerSocket(ctx, cfg.SocketPath, ms.Streamer(), mm.validateSegment(ctx), manifestSource) - recordWorkerExit("mkv", err, ctx.Err()) + recordWorkerExit("mp4", err, ctx.Err()) // Reap the worker unless we're deliberately leaving it running across a main // restart (ctx cancel). On a clean end OR a crash the worker has exited, so // Wait() clears the zombie; only on main shutdown do we let it stay detached diff --git a/pkg/media/ingest_daemon_test.go b/pkg/media/ingest_daemon_test.go index 18e36c4d..d683ce0e 100644 --- a/pkg/media/ingest_daemon_test.go +++ b/pkg/media/ingest_daemon_test.go @@ -59,7 +59,7 @@ func TestDetachedWorkerZeroDowntime(t *testing.T) { InputFD: 4, } - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) mediaR, mediaW, err := os.Pipe() require.NoError(t, err) @@ -67,7 +67,7 @@ func TestDetachedWorkerZeroDowntime(t *testing.T) { require.NoError(t, err) mediaR.Close() // the worker holds its own dup go func() { - _, _ = mediaW.Write(mkv) + _, _ = mediaW.Write(mp4) mediaW.Close() }() diff --git a/pkg/media/ingest_subprocess_test.go b/pkg/media/ingest_subprocess_test.go index 26a3d1d2..31c2115e 100644 --- a/pkg/media/ingest_subprocess_test.go +++ b/pkg/media/ingest_subprocess_test.go @@ -21,7 +21,7 @@ import ( // runIngestWorkerHelper is what the test binary becomes when re-exec'd with the // `ingest-worker` arg (see TestMain). It mirrors makeIngestWorkerCommand exactly: -// config on fd 3, frames on fd 4, MKV on stdin; clean run ends with End, a fatal +// config on fd 3, frames on fd 4, fMP4 on stdin; clean run ends with End, a fatal // error with an Error frame and a non-zero exit. func runIngestWorkerHelper() int { // Test hook: a worker-shaped process that just sleeps (same argv layout as a @@ -57,7 +57,7 @@ func runIngestWorkerHelper() int { defer f.Close() raw = f } - if err := ServeMKVIngestWorkerSocket(context.Background(), cfg, WorkerInput(cfg, raw)); err != nil { + if err := ServeMP4IngestWorkerSocket(context.Background(), cfg, WorkerInput(cfg, raw)); err != nil { return 1 } return 0 @@ -69,7 +69,7 @@ func runIngestWorkerHelper() int { } defer framesFile.Close() frames := ingestframe.NewWriter(framesFile) - if err := RunMKVIngestWorker(context.Background(), cfg, os.Stdin, frames, func() []byte { return cfg.Manifest }); err != nil { + if err := RunMP4IngestWorker(context.Background(), cfg, os.Stdin, frames, func() []byte { return cfg.Manifest }); err != nil { _ = frames.Error(err.Error()) return 1 } @@ -78,8 +78,8 @@ func runIngestWorkerHelper() int { } // TestIngestWorkerSubprocess exercises the real process boundary: it spawns the -// worker as an actual subprocess (config over fd 3, MKV over stdin, frames over -// fd 4 — the exact wiring MKVIngestIsolated uses) and verifies the worker +// worker as an actual subprocess (config over fd 3, fMP4 over stdin, frames over +// fd 4 — the exact wiring MP4IngestIsolated uses) and verifies the worker // produces valid signed segments, a clean End frame, and a zero exit. This is // the part the in-process worker test can't cover: fd passing, the framed wire // protocol over a real pipe, and process lifecycle. @@ -99,14 +99,14 @@ func TestIngestWorkerSubprocess(t *testing.T) { }) require.NoError(t, err) - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) exe, err := os.Executable() require.NoError(t, err) cmd := exec.CommandContext(ctx, exe, "ingest-worker") // Quiet gst in the child (it inherits the parent test's verbose leak-tracer env). cmd.Env = append(os.Environ(), "GST_DEBUG=0", "GST_TRACERS=") - cmd.Stdin = bytes.NewReader(mkv) + cmd.Stdin = bytes.NewReader(mp4) cmd.Stderr = os.Stderr cfgR, cfgW, err := os.Pipe() @@ -153,14 +153,14 @@ func TestIngestWorkerSubprocess(t *testing.T) { t.Logf("worker subprocess emitted %d valid signed segments + clean End", segs) } -// TestMKVIngestIsolatedWedgeContained is the isolation guarantee: an audio-only -// MKV starves the fMP4 muxer's video pad of both data and EOS, so the native +// TestMP4IngestIsolatedWedgeContained is the isolation guarantee: an audio-only +// fMP4 starves the fMP4 muxer's video pad of both data and EOS, so the native // pipeline wedges with no frames and 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 +// watchdog kills the worker and MP4IngestIsolated returns an error, bounded in // time, with THIS process — the node — still running to assert it. -func TestMKVIngestIsolatedWedgeContained(t *testing.T) { +func TestMP4IngestIsolatedWedgeContained(t *testing.T) { old := ingestWorkerWatchdog ingestWorkerWatchdog = 6 * time.Second defer func() { ingestWorkerWatchdog = old }() @@ -168,10 +168,10 @@ func TestMKVIngestIsolatedWedgeContained(t *testing.T) { mm, _ := getStaticTestMediaManager(t) ms := newBareSegmentSigner(t) - wedge := makeAudioOnlyAACMKV(t, context.Background(), 5) + wedge := makeAudioOnlyAACFMP4(t, context.Background(), 5) start := time.Now() - err := mm.MKVIngestIsolated(context.Background(), bytes.NewReader(wedge), ms) + err := mm.MP4IngestIsolated(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") @@ -208,7 +208,7 @@ func TestWorkerIngestsFromPassedFD(t *testing.T) { }) require.NoError(t, err) - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) exe, err := os.Executable() require.NoError(t, err) @@ -233,7 +233,7 @@ func TestWorkerIngestsFromPassedFD(t *testing.T) { cfgW.Close() }() go func() { - _, _ = mediaW.Write(mkv) + _, _ = mediaW.Write(mp4) mediaW.Close() }() diff --git a/pkg/media/ingest_supervisor.go b/pkg/media/ingest_supervisor.go index 5f09a984..d82091e8 100644 --- a/pkg/media/ingest_supervisor.go +++ b/pkg/media/ingest_supervisor.go @@ -27,7 +27,7 @@ import ( // shorten it. var ingestWorkerWatchdog = 30 * time.Second -// MKVIngestIsolated is the process-isolated counterpart to MKVIngest. Instead of +// MP4IngestIsolated is the process-isolated counterpart to MP4Ingest. 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 // `ingest-worker` subprocess that owns the pipeline and streams signed canonical @@ -38,7 +38,7 @@ var ingestWorkerWatchdog = 30 * time.Second // Per the locked design the worker signs everything, so main hands it the // streamer key + cert + a once-built manifest over a dedicated config fd (kept // off argv/env). See buildWorkerConfig for the interim key-custody note. -func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, ms MediaSigner) error { +func (mm *MediaManager) MP4IngestIsolated(ctx context.Context, input io.Reader, ms MediaSigner) error { cfg, err := mm.buildWorkerConfig(ctx, ms) if err != nil { return err @@ -104,7 +104,7 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, } cfgR.Close() // the child holds its own copy now framesW.Close() // ditto; the parent only reads framesR - spmetrics.IngestWorkerStarts.WithLabelValues("mkv-fd").Inc() + spmetrics.IngestWorkerStarts.WithLabelValues("mp4-fd").Inc() go func() { _, _ = cfgW.Write(cfgJSON) @@ -160,7 +160,7 @@ func (mm *MediaManager) MKVIngestIsolated(ctx context.Context, input io.Reader, } else if werr != nil && !sawEnd { exitErr = werr } - recordWorkerExit("mkv-fd", exitErr, ctx.Err()) + recordWorkerExit("mp4-fd", exitErr, ctx.Err()) switch { case readErr != nil: @@ -245,7 +245,7 @@ func streamWorkerLogs(ctx context.Context, stderr io.Reader, streamer string) { // buildWorkerConfig extracts the handshake the worker needs to sign on main's // behalf. INTERIM key custody: requires a software MediaSignerLocal — the -// MKV/RTMP push path always provides one; anything else errors so the caller can +// Mist-pull/RTMP path always provides one; anything else errors so the caller can // fall back to the in-process path. func (mm *MediaManager) buildWorkerConfig(ctx context.Context, ms MediaSigner) (IngestWorkerConfig, error) { local, ok := ms.(*MediaSignerLocal) diff --git a/pkg/media/ingest_worker.go b/pkg/media/ingest_worker.go index 8221738a..bdd17c42 100644 --- a/pkg/media/ingest_worker.go +++ b/pkg/media/ingest_worker.go @@ -50,7 +50,7 @@ func (h *manifestHolder) set(b []byte) { // approach; the detach/reattach work will revisit how a worker holds keys. type IngestWorkerConfig struct { StreamerDID string `json:"streamer_did"` - // KeyPEM is the streamer's ES256K signing key in PEM. The MKV/RTMP push path + // KeyPEM is the streamer's ES256K signing key in PEM. The Mist-pull/RTMP path // always yields a software key, which is what muxl-sign wants here; it is // forwarded verbatim, no reconstruction. KeyPEM []byte `json:"key_pem"` @@ -90,7 +90,7 @@ type IngestWorkerConfig struct { 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 + // (the fMP4 ingest 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 @@ -99,9 +99,9 @@ type IngestWorkerConfig struct { 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. + // Transport selects the worker's ingest source: "" / "mp4" reads fragmented + // MP4 media (stdin or InputFD); "whip" makes the worker own the WebRTC + // PeerConnection, built from OfferSDP — no media fd to pass. Transport string `json:"transport,omitempty"` // OfferSDP is the WHIP client's SDP offer (transport "whip"). The worker // generates the answer and emits it as the first frame (ingestframe.Answer) @@ -131,7 +131,7 @@ func WorkerInput(cfg IngestWorkerConfig, raw io.Reader) io.Reader { // the streamer key PEM + cert straight to muxl-sign, no MediaSigner / model / DB // needed. The manifest is read FRESH per GoP from the holder, so a manifest main // pushes mid-stream (e.g. pre-live → live) takes effect on the next GoP — the -// same fresh-per-GoP shape as the in-process signer. Shared by the MKV and WHIP +// same fresh-per-GoP shape as the in-process signer. Shared by the MP4 and WHIP // workers. func workerSignStream(cfg IngestWorkerConfig, getManifest func() []byte) SignSegmentStreamFunc { return func(ctx context.Context, input io.Reader, eventCh chan *muxl.MuxlEvent) error { @@ -152,7 +152,7 @@ func workerSignStream(cfg IngestWorkerConfig, getManifest func() []byte) SignSeg // segment); flush Closes that transcoder so its ~1-GoP tail is framed before the // worker exits. The transcoder runs on a non-cancellable context so draining the // signer can't kill it early. One process == one session, so the per-DID -// transcoder-reuse hazard can't arise. Shared by the MKV and WHIP workers. +// transcoder-reuse hazard can't arise. Shared by the MP4 and WHIP workers. func (mm *MediaManager) workerSegmentSink(ctx context.Context, cfg IngestWorkerConfig, frames FrameWriter) (onSegment func(context.Context, []byte) error, flush func()) { var transcoder *streamTranscoder onSegment = func(_ context.Context, segment []byte) error { @@ -187,9 +187,9 @@ func (mm *MediaManager) workerSegmentSink(ctx context.Context, cfg IngestWorkerC return onSegment, flush } -// RunMKVIngestWorker is the body of the `ingest-worker` subcommand. It reads an -// MKV stream from stdin, runs the same demux + Opus re-encode + muxl-sign -// pipeline as the in-process MKVIngest, and emits each signed canonical .m4s +// RunMP4IngestWorker is the body of the `ingest-worker` subcommand. It reads a +// fragmented-MP4 stream from stdin, runs the same demux + Opus re-encode + muxl-sign +// pipeline as the in-process MP4Ingest, and emits each signed canonical .m4s // segment to frames; the main process reads those frames and runs ValidateMP4 // over each, exactly as if onSegment had called it directly. // @@ -197,7 +197,7 @@ func (mm *MediaManager) workerSegmentSink(ctx context.Context, cfg IngestWorkerC // caller frames End or Error accordingly. All segment frames are guaranteed // flushed before it returns, so a trailing End can never race ahead of the last // Segment. -func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader, frames FrameWriter, getManifest func() []byte) error { +func RunMP4IngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Reader, frames FrameWriter, getManifest func() []byte) error { gstinit.InitGST() ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -224,7 +224,7 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re pr, pw := io.Pipe() media = io.TeeReader(stdin, pw) go func() { - if derr := mm.dumpToFile(ctx, pr, cfg.StreamerDID, ".rtmp.mkv"); derr != nil { + if derr := mm.dumpToFile(ctx, pr, cfg.StreamerDID, ".rtmp.mp4"); derr != nil { log.Error(ctx, "ingest worker: dump recording to file", "error", derr) } }() @@ -234,7 +234,7 @@ func RunMKVIngestWorker(ctx context.Context, cfg IngestWorkerConfig, stdin io.Re if err != nil { return fmt.Errorf("build signer element: %w", err) } - pipeline, err := buildMKVIngestPipeline(ctx, media, signerElem) + pipeline, err := buildMP4IngestPipeline(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 2615a631..ee1bd5b3 100644 --- a/pkg/media/ingest_worker_test.go +++ b/pkg/media/ingest_worker_test.go @@ -25,7 +25,7 @@ import ( // fd-passed push connection: prebuf (bytes main read past the headers) is // prepended, and a chunked transfer-encoding is decoded back to the raw media. func TestWorkerInputDeframes(t *testing.T) { - payload := []byte("the-actual-media-bytes-pretend-this-is-mkv-data") + payload := []byte("the-actual-media-bytes-pretend-this-is-mp4-data") // A textbook chunked body: one chunk then the zero terminator. body := []byte(fmt.Sprintf("%x\r\n%s\r\n0\r\n\r\n", len(payload), payload)) @@ -41,18 +41,15 @@ func TestWorkerInputDeframes(t *testing.T) { require.Equal(t, payload, got) } -// makeH264AACMKV builds a clean, single-track, streamable H264+AAC MKV from an -// H264+Opus MP4 fixture (video passed through, audio transcoded Opus→AAC). The -// repo's only AAC fixture (sample-stream.mkv) carries four audio tracks, which -// the single-audio ingest pipeline leaves three of unlinked — wedging -// matroskademux with no EOS. This produces exactly the 1-video-1-audio AAC MKV -// the RTMP push path actually delivers. -func makeH264AACMKV(t *testing.T, ctx context.Context, srcMP4 string) []byte { +// makeH264AACFMP4 builds a clean, 1-video-1-audio fragmented H264+AAC MP4 from +// an H264+Opus MP4 fixture (video passed through, audio transcoded Opus→AAC) — +// exactly the shape MistServer's live .mp4 output delivers for an RTMP push. +func makeH264AACFMP4(t *testing.T, ctx context.Context, srcMP4 string) []byte { t.Helper() gstinit.InitGST() desc := strings.Join([]string{ "filesrc location=" + srcMP4 + " ! qtdemux name=d", - "d. ! queue ! h264parse ! matroskamux name=mux streamable=true ! appsink name=sink", + "d. ! queue ! h264parse ! mp4mux name=mux fragment-duration=500 ! appsink name=sink", "d. ! queue ! opusdec ! audioconvert ! audioresample ! fdkaacenc ! aacparse ! mux.", }, "\n") pipeline, err := gst.NewPipelineFromString(desc) @@ -69,23 +66,21 @@ func makeH264AACMKV(t *testing.T, ctx context.Context, srcMP4 string) []byte { go func() { busErr <- HandleBusMessages(ctx, pipeline) }() require.NoError(t, pipeline.SetState(gst.StatePlaying)) defer func() { _ = pipeline.SetState(gst.StateNull) }() - require.NoError(t, <-busErr, "remux to H264+AAC MKV") - require.NotEmpty(t, buf.Bytes(), "remux produced an MKV") + require.NoError(t, <-busErr, "remux to fragmented H264+AAC MP4") + require.NotEmpty(t, buf.Bytes(), "remux produced an fMP4") return buf.Bytes() } -// makeAudioOnlyAACMKV synthesizes an AAC-audio-only streamable MKV — the +// makeAudioOnlyAACFMP4 synthesizes an AAC-audio-only fragmented MP4 — the // canonical WEDGE input for watchdog/containment tests. The ingest pipeline -// hardwires a video and an audio branch; with no video track, matroskademux -// never creates a video pad, the fMP4 aggregator's video pad never sees data -// OR EOS, and the pipeline hangs forever with no frames and no EOS — a true -// native wedge that no queue sizing can fix. (The 4-audio sample-stream.mkv -// previously used for this stopped wedging once the ingest branches moved to -// Queue2Big: its wedge was really the 1s default-queue interleave deadlock.) -func makeAudioOnlyAACMKV(t *testing.T, ctx context.Context, seconds int) []byte { +// hardwires a video and an audio branch; with no video track, qtdemux never +// creates a video pad, the fMP4 aggregator's video pad never sees data OR +// EOS, and the pipeline hangs forever with no frames and no EOS — a true +// native wedge that no queue sizing can fix. +func makeAudioOnlyAACFMP4(t *testing.T, ctx context.Context, seconds int) []byte { t.Helper() gstinit.InitGST() - desc := fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! matroskamux streamable=true ! appsink name=sink", seconds*47) + desc := fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! mp4mux fragment-duration=500 ! appsink name=sink", seconds*47) pipeline, err := gst.NewPipelineFromString(desc) require.NoError(t, err) @@ -100,18 +95,18 @@ func makeAudioOnlyAACMKV(t *testing.T, ctx context.Context, seconds int) []byte go func() { busErr <- HandleBusMessages(ctx, pipeline) }() require.NoError(t, pipeline.SetState(gst.StatePlaying)) defer func() { _ = pipeline.SetState(gst.StateNull) }() - require.NoError(t, <-busErr, "synthesize audio-only MKV") + require.NoError(t, <-busErr, "synthesize audio-only fMP4") require.NotEmpty(t, buf.Bytes()) return buf.Bytes() } -// TestRunMKVIngestWorkerProducesValidSignedFrames drives the isolated ingest -// worker's core directly (no subprocess): feed it an H264+AAC MKV, collect the +// TestRunMP4IngestWorkerProducesValidSignedFrames drives the isolated ingest +// worker's core directly (no subprocess): feed it an H264+AAC fMP4, collect the // framed output, and verify every emitted segment is a valid signed canonical // .m4s. This is the contract the supervisor relies on — frames it can hand // straight to ValidateMP4. (The real subprocess spawn + fault injection is // Stage 3.) -func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { +func TestRunMP4IngestWorkerProducesValidSignedFrames(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) @@ -132,13 +127,13 @@ func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { BroadcasterHost: "test.example.com", } - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) - // All frame writes complete before RunMKVIngestWorker returns (it waits on + // All frame writes complete before RunMP4IngestWorker returns (it waits on // the signer drain), so reading the buffer single-threaded afterwards is safe. var buf bytes.Buffer frames := ingestframe.NewWriter(&buf) - require.NoError(t, RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), frames, func() []byte { return cfg.Manifest })) + require.NoError(t, RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), frames, func() []byte { return cfg.Manifest })) r := ingestframe.NewReader(&buf) var segs int @@ -175,12 +170,12 @@ func TestRunMKVIngestWorkerProducesValidSignedFrames(t *testing.T) { t.Logf("worker emitted %d valid dual-codec segments", segs) } -// TestRunMKVIngestWorkerRecords proves debug recording works INSIDE the worker: +// TestRunMP4IngestWorkerRecords 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 +// media to debug-recordings//.rtmp.mp4. 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) { +func TestRunMP4IngestWorkerRecords(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) @@ -200,33 +195,33 @@ func TestRunMKVIngestWorkerRecords(t *testing.T) { DataDir: dataDir, } - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(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), func() []byte { return cfg.Manifest })) + require.NoError(t, RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest })) - // The recording lands at debug-recordings//.rtmp.mkv and + // The recording lands at debug-recordings//.rtmp.mp4 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") + glob := filepath.Join(dataDir, "debug-recordings", "*", "*.rtmp.mp4") 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) + return rerr == nil && bytes.Equal(got, mp4) }, 10*time.Second, 25*time.Millisecond, "worker records the ingest media verbatim") } -// TestRunMKVIngestWorkerSelfWatchdog proves the worker's OWN watchdog contains a +// TestRunMP4IngestWorkerSelfWatchdog 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. An audio-only MKV starves the muxer's video pad of both data and +// exit itself. An audio-only fMP4 starves the muxer's video pad of both data and // EOS, so the pipeline wedges with no frames; the watchdog must tear it 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) { +// watchdog is covered separately by TestMP4IngestIsolatedWedgeContained.) +func TestRunMP4IngestWorkerSelfWatchdog(t *testing.T) { old := ingestWorkerWatchdog ingestWorkerWatchdog = 3 * time.Second defer func() { ingestWorkerWatchdog = old }() @@ -245,12 +240,12 @@ func TestRunMKVIngestWorkerSelfWatchdog(t *testing.T) { BroadcasterHost: "test.example.com", } - wedge := makeAudioOnlyAACMKV(t, ctx, 5) + wedge := makeAudioOnlyAACFMP4(t, ctx, 5) start := time.Now() done := make(chan error, 1) go func() { - done <- RunMKVIngestWorker(ctx, cfg, bytes.NewReader(wedge), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest }) + done <- RunMP4IngestWorker(ctx, cfg, bytes.NewReader(wedge), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest }) }() select { case <-done: @@ -258,16 +253,16 @@ func TestRunMKVIngestWorkerSelfWatchdog(t *testing.T) { 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)") + t.Fatal("worker-side watchdog did not contain the wedge (RunMP4IngestWorker hung)") } } -// TestRunMKVIngestWorkerSignsWithSuppliedManifest is the core of the pre-live → +// TestRunMP4IngestWorkerSignsWithSuppliedManifest is the core of the pre-live → // live fix: the worker signs each GoP with whatever the manifest getter returns, // NOT a frozen one. Two runs of the same media differ only in the getter — a // pre-live manifest yields unpublished segments; a live manifest (the same plus // a c2pa.published action, as main would push on go-live) yields published ones. -func TestRunMKVIngestWorkerSignsWithSuppliedManifest(t *testing.T) { +func TestRunMP4IngestWorkerSignsWithSuppliedManifest(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) @@ -285,11 +280,11 @@ func TestRunMKVIngestWorkerSignsWithSuppliedManifest(t *testing.T) { []byte(`{"action":"c2pa.created"},{"action":"c2pa.published"}`), 1) require.NotEqual(t, string(prelive), string(live), "the live manifest must add c2pa.published") - mkv := makeH264AACMKV(t, ctx, getFixture("5sec.mp4")) + mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) firstSegmentPublished := func(manifest []byte) bool { var buf bytes.Buffer - require.NoError(t, RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), + require.NoError(t, RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), ingestframe.NewWriter(&buf), func() []byte { return manifest })) r := ingestframe.NewReader(&buf) typ, payload, rerr := r.ReadFrame() diff --git a/pkg/media/key_revocation_test.go b/pkg/media/key_revocation_test.go index 065d2baf..35cfd91c 100644 --- a/pkg/media/key_revocation_test.go +++ b/pkg/media/key_revocation_test.go @@ -87,20 +87,20 @@ func TestWatchKeyRevocationStreamKick(t *testing.T) { } } -// TestMKVIngestIsolatedBanContained proves the fix end to end: banning a streamer +// TestMP4IngestIsolatedBanContained proves the fix end to end: banning a streamer // mid-ingest tears their isolated worker down. The watchdog is set generously -// (60s) and the input is a wedging audio-only MKV that never ends on its own — +// (60s) and the input is a wedging audio-only fMP4 that never ends on its own — // so a timely return can only be the ban kill, not the watchdog or a natural -// EOS. (The 4-audio sample-stream.mkv previously used here now ingests to +// EOS. (The 4-audio sample-stream.mp4 previously used here now ingests to // completion in a few seconds, which would race the ban.) -func TestMKVIngestIsolatedBanContained(t *testing.T) { +func TestMP4IngestIsolatedBanContained(t *testing.T) { old := ingestWorkerWatchdog ingestWorkerWatchdog = 60 * time.Second defer func() { ingestWorkerWatchdog = old }() mm, _ := getStaticTestMediaManager(t) ms := newBareSegmentSigner(t) - wedge := makeAudioOnlyAACMKV(t, context.Background(), 5) + wedge := makeAudioOnlyAACFMP4(t, context.Background(), 5) // Ban the streamer once the worker is up and the watcher has subscribed. go func() { @@ -112,7 +112,7 @@ func TestMKVIngestIsolatedBanContained(t *testing.T) { }() start := time.Now() - err := mm.MKVIngestIsolated(context.Background(), bytes.NewReader(wedge), ms) + err := mm.MP4IngestIsolated(context.Background(), bytes.NewReader(wedge), ms) elapsed := time.Since(start) require.Error(t, err, "a banned stream is torn down, surfaced as an error") diff --git a/pkg/media/mist_mkv_ingest_test.go b/pkg/media/mist_mkv_ingest_test.go deleted file mode 100644 index 6b67609d..00000000 --- a/pkg/media/mist_mkv_ingest_test.go +++ /dev/null @@ -1,274 +0,0 @@ -package media - -import ( - "bytes" - "context" - "errors" - "fmt" - "io" - "os" - "strings" - "testing" - "time" - - "github.com/go-gst/go-gst/gst" - "github.com/go-gst/go-gst/gst/app" - "github.com/stretchr/testify/require" - "stream.place/streamplace/pkg/crypto/signers" - "stream.place/streamplace/pkg/gstinit" - "stream.place/streamplace/pkg/ingestframe" - "stream.place/streamplace/test/remote" -) - -// runMKVThroughIngestWorker feeds an MKV byte stream through the isolated -// ingest worker and returns how many signed canonical segments it emitted. The -// worker watchdog is shortened so a wedged pipeline returns promptly instead of -// hanging the test (override via SP_TEST_WATCHDOG to e.g. park a wedge for a -// stack dump). With transcode=true the node keys are supplied so the worker -// completes to dual-codec, as production does. -func runMKVThroughIngestWorker(t *testing.T, mkv []byte, transcode bool) (int, error) { - segs, err := runMKVThroughIngestWorkerSegments(t, mkv, transcode) - return len(segs), err -} - -// runMKVThroughIngestWorkerSegments is runMKVThroughIngestWorker returning the -// raw segment payloads, for tests that validate the emitted segments rather -// than just count them. -func runMKVThroughIngestWorkerSegments(t *testing.T, mkv []byte, transcode bool) ([][]byte, error) { - t.Helper() - old := ingestWorkerWatchdog - ingestWorkerWatchdog = 10 * time.Second - if wd := os.Getenv("SP_TEST_WATCHDOG"); wd != "" { - d, perr := time.ParseDuration(wd) - require.NoError(t, perr) - ingestWorkerWatchdog = d - } - 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", - } - if transcode { - cfg.NodeCertPEM = ms.Cert - cfg.NodeKeyPEM = keyPEM - } - - var buf bytes.Buffer - runErr := RunMKVIngestWorker(ctx, cfg, bytes.NewReader(mkv), ingestframe.NewWriter(&buf), func() []byte { return cfg.Manifest }) - - r := ingestframe.NewReader(&buf) - var segs [][]byte - for { - typ, payload, rerr := r.ReadFrame() - if errors.Is(rerr, io.EOF) { - break - } - require.NoError(t, rerr) - if typ == ingestframe.Segment { - segs = append(segs, payload) - } - } - return segs, runErr -} - -// makeSparseVideoAACMKV synthesizes the stream shape that wedged production -// ingest: video that degrades to keyframe-only at a low rate (here 0.5fps — -// MistServer drops all delta frames when a push falls behind, leaving ~1s-apart -// keyframes) alongside continuous 48kHz AAC audio, in a streamable MKV. The -// audio branch must buffer a full video-frame gap while matroskademux walks to -// the next video frame; gst's default 1s-capped queue can't, and the -// aggregator-based fMP4 muxer deadlocks (see buildMKVIngestPipeline). -func makeSparseVideoAACMKV(t *testing.T, ctx context.Context, seconds int) []byte { - t.Helper() - gstinit.InitGST() - desc := strings.Join([]string{ - fmt.Sprintf("videotestsrc num-buffers=%d ! video/x-raw,width=320,height=240,framerate=1/2 ! x264enc key-int-max=1 tune=zerolatency speed-preset=ultrafast ! h264parse ! matroskamux name=mux streamable=true ! appsink name=sink", (seconds+1)/2), - fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! mux.", seconds*47), - }, "\n") - pipeline, err := gst.NewPipelineFromString(desc) - require.NoError(t, err) - - sinkEle, err := pipeline.GetElementByName("sink") - require.NoError(t, err) - var buf bytes.Buffer - app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, &buf), - }) - - busErr := make(chan error, 1) - go func() { busErr <- HandleBusMessages(ctx, pipeline) }() - require.NoError(t, pipeline.SetState(gst.StatePlaying)) - defer func() { _ = pipeline.SetState(gst.StateNull) }() - require.NoError(t, <-busErr, "synthesize sparse-video MKV") - require.NotEmpty(t, buf.Bytes()) - return buf.Bytes() -} - -// TestMKVIngestSparseVideoNoWedge is the regression test for a production -// ingest wedge: a stream whose video goes sparse (keyframe-only, ≥1s between -// video frames) deadlocked the MKV ingest pipeline — audio backpressure -// through the default 1s-capped queue blocked the demux, the fMP4 aggregator -// starved on its video pad, and the stream hung with no EOS until the -// watchdog killed it. With Queue2Big on both ingest branches the same stream -// must segment to completion. -func TestMKVIngestSparseVideoNoWedge(t *testing.T) { - ctx := context.Background() - mkv := makeSparseVideoAACMKV(t, ctx, 12) - - segs, err := runMKVThroughIngestWorker(t, mkv, false) - require.NoError(t, err, "sparse-video stream ingests cleanly (wedge → watchdog → context canceled)") - // 12s of 2s-apart keyframes ≈ 6 GoPs; wedging yields 0-1 segments. - require.GreaterOrEqual(t, segs, 4, "sparse-video stream emits its segments") - t.Logf("sparse-video stream: %d segments", segs) -} - -// The nyc-* fixtures are cuts of a real 164s production MistServer MKV push -// whose video degrades to keyframe-only at ~140s (behind-push frame-drop) — -// the capture that wedged production ingest. head = first ~10s; head-nojson = -// the same bytes with the 30-byte M_JSON TrackEntry stripped; tail135 = from -// 135s (~5s before the keyframe-only transition); full = the whole capture. - -// TestMKVIngestMistMetadataTrack: MistServer's MKV push declares a -// live-metadata track (CodecID M_JSON, TrackType 3) as track 1, ahead of the -// AAC audio and H264 video tracks. It was the initial suspect for the -// production wedge but proved benign — matroskademux ignores the unknown -// codec, and the same media segments identically with the 30-byte M_JSON -// TrackEntry stripped (the control). Kept as a canary for the MistServer -// track layout. (The real wedge: TestMKVIngestSparseVideoNoWedge.) -func TestMKVIngestMistMetadataTrack(t *testing.T) { - control, err := os.ReadFile(remote.RemoteFixture("3284ef5658e7864bce326c296a909e985c4167d0b9a445b2ce944c2f0171c71e/nyc-head-nojson.mkv")) - require.NoError(t, err) - mist, err := os.ReadFile(remote.RemoteFixture("c0989e044f3350c55f1e129b76252bfb2859914058bb17d2431a605db9693467/nyc-head.mkv")) - require.NoError(t, err) - - segs, err := runMKVThroughIngestWorker(t, control, false) - require.NoError(t, err, "control (M_JSON TrackEntry stripped) ingests cleanly") - require.GreaterOrEqual(t, segs, 1, "control emits segments") - t.Logf("control: %d segments", segs) - - segs, err = runMKVThroughIngestWorker(t, mist, false) - require.NoError(t, err, "MistServer MKV (with M_JSON track) ingests cleanly") - require.GreaterOrEqual(t, segs, 1, "MistServer MKV emits segments") - t.Logf("with M_JSON track: %d segments", segs) -} - -// TestMKVIngestMistFullSample runs the entire 164s production capture through -// the worker with node transcode keys — the closest in-process approximation -// of the production ingest. The capture degrades to keyframe-only video at -// ~140s (MistServer behind-push frame-drop), which is what wedged production; -// with Queue2Big on the ingest branches the whole capture must segment. -func TestMKVIngestMistFullSample(t *testing.T) { - // This once "progressively slowed until the watchdog fired, then hung in - // the post-cancel drain" and was skip-gated as known-hanging — that was - // the muxl-event-drain-vs-cancel deadlock (see muxlSignSegmentElem's - // drainCtx); with the drain non-cancellable the full capture transcodes - // at full speed (~13s). - mkv, err := os.ReadFile(remote.RemoteFixture("3e4e5d9758e67053908e523379a3e2ef2cf60679d0657a940daf96590e866015/nyc-full.mkv")) - require.NoError(t, err) - segs, err := runMKVThroughIngestWorker(t, mkv, true) - t.Logf("full sample: %d segments, err=%v", segs, err) - require.NoError(t, err, "full production sample ingests cleanly") - // 171 GoPs in the capture (~1s each); wedging yielded ~144. - require.GreaterOrEqual(t, segs, 160, "full sample emits the whole stream's segments") -} - -// TestMKVIngestMistTail is the fast sample-based wedge check: the tail sample -// starts at 135s, ~5s before the capture goes keyframe-only, so an unfixed -// pipeline wedges within seconds (4 segments) instead of minutes. -func TestMKVIngestMistTail(t *testing.T) { - mkv, err := os.ReadFile(remote.RemoteFixture("03df698a342f1ab89dccc20ce0a0283e1270104e6382703575686c9f4a88881e/nyc-tail135.mkv")) - require.NoError(t, err) - segs, err := runMKVThroughIngestWorker(t, mkv, false) - t.Logf("tail sample: %d segments, err=%v", segs, err) - require.NoError(t, err, "tail of production sample ingests cleanly") - // 135s..164s at ~1s GoPs ≈ 29 segments; wedging yields ~5. - require.GreaterOrEqual(t, segs, 20, "tail sample emits segments past the keyframe-only transition") -} - -// makeBFrameAACMKV synthesizes an H264 stream WITH B-frames (PTS ≠ DTS) -// alongside AAC audio in a streamable MKV — the shape a hardware encoder or -// non-zerolatency x264 push produces. Matroska blocks carry only presentation -// timestamps, so on demux the reordered video arrives with dts=none; without -// DTS reconstruction the fMP4 muxer stretches the video track and every -// segment fails validation downstream (see buildMKVIngestPipeline's -// h264timestamper). b-adapt=false forces x264 to actually emit the configured -// B-frames rather than deciding per-scene. -func makeBFrameAACMKV(t *testing.T, ctx context.Context, seconds int) []byte { - t.Helper() - gstinit.InitGST() - desc := strings.Join([]string{ - fmt.Sprintf("videotestsrc num-buffers=%d ! video/x-raw,width=320,height=240,framerate=30/1 ! x264enc bframes=2 b-adapt=false key-int-max=30 speed-preset=veryfast ! h264parse ! matroskamux name=mux streamable=true ! appsink name=sink", seconds*30), - fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! mux.", seconds*47), - }, "\n") - pipeline, err := gst.NewPipelineFromString(desc) - require.NoError(t, err) - - sinkEle, err := pipeline.GetElementByName("sink") - require.NoError(t, err) - var buf bytes.Buffer - app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, &buf), - }) - - busErr := make(chan error, 1) - go func() { busErr <- HandleBusMessages(ctx, pipeline) }() - require.NoError(t, pipeline.SetState(gst.StatePlaying)) - defer func() { _ = pipeline.SetState(gst.StateNull) }() - require.NoError(t, <-busErr, "synthesize B-frame MKV") - require.NotEmpty(t, buf.Bytes()) - return buf.Bytes() -} - -// TestMKVIngestBFramesValidate is the regression test for B-frame MKV ingest: -// every emitted segment must survive the full ValidateMP4Media chokepoint. -// Before DTS reconstruction, mp4mux treated the reordered (dts=none) B-frame -// PTS as monotonic timing, stretching the video track ~2.2×; the segments -// LOOKED fine (both tracks present, signatures valid) but the video/audio -// duration mismatch made push-mode qtdemux EOS the audio pad before the audio -// bytes arrived — muxl's flat wrap is non-interleaved, video first — and every -// segment was rejected with "no audio in segment". -func TestMKVIngestBFramesValidate(t *testing.T) { - ctx := context.Background() - mkv := makeBFrameAACMKV(t, ctx, 8) - - segs, err := runMKVThroughIngestWorkerSegments(t, mkv, false) - require.NoError(t, err, "B-frame stream ingests cleanly") - require.GreaterOrEqual(t, len(segs), 6, "B-frame stream emits its segments") - - sawBFrames := false - for i, seg := range segs { - res, verr := ValidateMP4Media(ctx, seg) - require.NoError(t, verr, "segment %d validates (video+audio both present)", i) - dur := time.Duration(res.MediaData.Duration) - require.Greater(t, dur, 500*time.Millisecond, "segment %d duration sane", i) - require.Less(t, dur, 2*time.Second, "segment %d duration not stretched", i) - if res.MediaData.Video[0].BFrames { - sawBFrames = true - } - } - require.True(t, sawBFrames, "synthesized stream actually contains B-frames — if this fails the test no longer exercises the reorder path") - t.Logf("B-frame stream: %d segments, all validated", len(segs)) -} - -// TestMKVIngestMistFullSampleNoTranscode is TestMKVIngestMistFullSample -// without node keys (segment+sign only). During diagnosis this proved the -// wedge was in the core ingest pipeline, not the dual-codec transcode stage — -// both variants wedged at the same GoP. -func TestMKVIngestMistFullSampleNoTranscode(t *testing.T) { - mkv, err := os.ReadFile(remote.RemoteFixture("3e4e5d9758e67053908e523379a3e2ef2cf60679d0657a940daf96590e866015/nyc-full.mkv")) - require.NoError(t, err) - segs, err := runMKVThroughIngestWorker(t, mkv, false) - t.Logf("full sample (no transcode): %d segments, err=%v", segs, err) - require.NoError(t, err, "full production sample ingests cleanly without transcode") - require.GreaterOrEqual(t, segs, 70, "full sample emits the whole stream's segments") -} diff --git a/pkg/media/mist_mp4_ingest_test.go b/pkg/media/mist_mp4_ingest_test.go new file mode 100644 index 00000000..13a60bcc --- /dev/null +++ b/pkg/media/mist_mp4_ingest_test.go @@ -0,0 +1,284 @@ +package media + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "os" + "strings" + "testing" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/pkg/gstinit" + "stream.place/streamplace/pkg/ingestframe" + "stream.place/streamplace/pkg/muxl" + "stream.place/streamplace/test/remote" +) + +// runMP4ThroughIngestWorker feeds a fragmented-MP4 byte stream through the +// isolated ingest worker and returns how many signed canonical segments it +// emitted. The worker watchdog is shortened so a wedged pipeline returns +// promptly instead of hanging the test (override via SP_TEST_WATCHDOG to e.g. +// park a wedge for a stack dump). With transcode=true the node keys are +// supplied so the worker completes to dual-codec, as production does. +func runMP4ThroughIngestWorker(t *testing.T, mp4 []byte, transcode bool) (int, error) { + segs, err := runMP4ThroughIngestWorkerSegments(t, mp4, transcode) + return len(segs), err +} + +// runMP4ThroughIngestWorkerSegments is runMP4ThroughIngestWorker returning the +// raw segment payloads, for tests that validate the emitted segments rather +// than just count them. +func runMP4ThroughIngestWorkerSegments(t *testing.T, mp4 []byte, transcode bool) ([][]byte, error) { + t.Helper() + old := ingestWorkerWatchdog + ingestWorkerWatchdog = 10 * time.Second + if wd := os.Getenv("SP_TEST_WATCHDOG"); wd != "" { + d, perr := time.ParseDuration(wd) + require.NoError(t, perr) + ingestWorkerWatchdog = d + } + 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", + } + if transcode { + cfg.NodeCertPEM = ms.Cert + cfg.NodeKeyPEM = keyPEM + } + + var buf bytes.Buffer + runErr := RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), ingestframe.NewWriter(&buf), func() []byte { return cfg.Manifest }) + + r := ingestframe.NewReader(&buf) + var segs [][]byte + for { + typ, payload, rerr := r.ReadFrame() + if errors.Is(rerr, io.EOF) { + break + } + require.NoError(t, rerr) + if typ == ingestframe.Segment { + segs = append(segs, payload) + } + } + return segs, runErr +} + +// runSynthPipeline runs a gst-launch description whose sink is an appsink +// named "sink" and returns everything the sink produced. +func runSynthPipeline(t *testing.T, ctx context.Context, desc string) []byte { + t.Helper() + gstinit.InitGST() + pipeline, err := gst.NewPipelineFromString(desc) + require.NoError(t, err) + + sinkEle, err := pipeline.GetElementByName("sink") + require.NoError(t, err) + var buf bytes.Buffer + app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: WriterNewSample(ctx, &buf), + }) + + busErr := make(chan error, 1) + go func() { busErr <- HandleBusMessages(ctx, pipeline) }() + require.NoError(t, pipeline.SetState(gst.StatePlaying)) + defer func() { _ = pipeline.SetState(gst.StateNull) }() + require.NoError(t, <-busErr, "synthesize test stream") + require.NotEmpty(t, buf.Bytes()) + return buf.Bytes() +} + +// makeSparseVideoAACFMP4 synthesizes the stream shape that wedged production +// ingest back when the bridge format was MKV: video that degrades to +// keyframe-only at a low rate (here 0.5fps — MistServer drops all delta frames +// when a push falls behind, leaving ~1s-apart keyframes) alongside continuous +// 48kHz AAC audio, in a fragmented MP4. The wedge mechanism is +// demux-agnostic — the audio branch must buffer a full video-frame gap while +// the demux walks the byte stream to the next video frame, and the +// aggregator-based fMP4 muxer downstream consumes nothing until every pad has +// data — so the regression carries over to the qtdemux pipeline (see +// buildMP4IngestPipeline's Queue2Big comment). +func makeSparseVideoAACFMP4(t *testing.T, ctx context.Context, seconds int) []byte { + t.Helper() + desc := strings.Join([]string{ + fmt.Sprintf("videotestsrc num-buffers=%d ! video/x-raw,width=320,height=240,framerate=1/2 ! x264enc key-int-max=1 tune=zerolatency speed-preset=ultrafast ! h264parse ! mp4mux name=mux fragment-duration=500 ! appsink name=sink", (seconds+1)/2), + fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! mux.", seconds*47), + }, "\n") + return runSynthPipeline(t, ctx, desc) +} + +// TestMP4IngestSparseVideoNoWedge is the regression test for a production +// ingest wedge: a stream whose video goes sparse (keyframe-only, ≥1s between +// video frames) deadlocked the ingest pipeline — audio backpressure through +// the default 1s-capped queue blocked the demux, the fMP4 aggregator starved +// on its video pad, and the stream hung with no EOS until the watchdog killed +// it. With Queue2Big on both ingest branches the same stream must segment to +// completion. +func TestMP4IngestSparseVideoNoWedge(t *testing.T) { + ctx := context.Background() + mp4 := makeSparseVideoAACFMP4(t, ctx, 12) + + segs, err := runMP4ThroughIngestWorker(t, mp4, false) + require.NoError(t, err, "sparse-video stream ingests cleanly (wedge → watchdog → context canceled)") + // 12s of 2s-apart keyframes ≈ 6 GoPs; wedging yields 0-1 segments. + require.GreaterOrEqual(t, segs, 4, "sparse-video stream emits its segments") + t.Logf("sparse-video stream: %d segments", segs) +} + +// videoPTSDTSOffsets flat-wraps a canonical segment and returns every video +// sample's PTS−DTS composition offset. This is the direct probe for the class +// of bug that motivated the fMP4 ingest rewrite: MKV ingest reconstructed DTS +// with h264timestamper, whose SPS fallback minted a constant spurious offset +// for streams that declare no reorder window (VideoToolbox) — pushing every +// GoP's presentation past its segment's declared window and breaking WebRTC +// playback at each keyframe. fMP4 ingest reads the container's real DTS, so a +// no-reorder stream must come out with PTS == DTS on every sample. +func videoPTSDTSOffsets(t *testing.T, ctx context.Context, segment []byte) []time.Duration { + t.Helper() + gstinit.InitGST() + var flat bytes.Buffer + require.NoError(t, muxl.RunMuxlWrap(ctx, bytes.NewReader(segment), "flat", &flat)) + + desc := strings.Join([]string{ + "appsrc name=src ! qtdemux name=demux", + "demux.video_0 ! queue ! h264parse ! appsink name=sink sync=false", + "demux.audio_0 ! queue ! fakesink sync=false", + }, "\n") + pipeline, err := gst.NewPipelineFromString(desc) + require.NoError(t, err) + + srcEle, err := pipeline.GetElementByName("src") + require.NoError(t, err) + app.SrcFromElement(srcEle).SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: ReaderNeedDataIncremental(ctx, bytes.NewReader(flat.Bytes())), + }) + + sinkEle, err := pipeline.GetElementByName("sink") + require.NoError(t, err) + var offsets []time.Duration + app.SinkFromElement(sinkEle).SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowEOS + } + buf := sample.GetBuffer() + pts, dts := buf.PresentationTimestamp(), buf.DecodingTimestamp() + if pts != gst.ClockTimeNone && dts != gst.ClockTimeNone { + offsets = append(offsets, time.Duration(int64(pts)-int64(dts))) + } + return gst.FlowOK + }, + }) + + busErr := make(chan error, 1) + go func() { busErr <- HandleBusMessages(ctx, pipeline) }() + require.NoError(t, pipeline.SetState(gst.StatePlaying)) + defer func() { _ = pipeline.SetState(gst.StateNull) }() + require.NoError(t, <-busErr, "demux flat-wrapped segment") + require.NotEmpty(t, offsets, "segment has video samples with timestamps") + return offsets +} + +// TestMP4IngestMistRealSample runs a real MistServer live-MP4 capture through +// the worker: a VideoToolbox (macOS hardware encoder) 720p H264 + AAC RTMP +// push, pulled from Mist's HTTP .mp4 output — exactly what production ingest +// consumes since the MKV→fMP4 bridge rewrite. VideoToolbox is the interesting +// encoder here because its SPS declares no reorder window, which is what sent +// the old MKV path's h264timestamper into its spurious-offset fallback; this +// stream must instead come through with its real timestamps: PTS == DTS on +// every video sample of every signed segment. +func TestMP4IngestMistRealSample(t *testing.T) { + ctx := context.Background() + mp4, err := os.ReadFile(remote.RemoteFixture("ee4d7f8f9b267ba8229314162ef268186048f91ac4242b13fce3f5ee955b97ae/mist-vt-720p.mp4")) + require.NoError(t, err) + + segs, err := runMP4ThroughIngestWorkerSegments(t, mp4, false) + require.NoError(t, err, "real Mist fMP4 capture ingests cleanly") + require.GreaterOrEqual(t, len(segs), 3, "capture emits its segments") + + for i, seg := range segs { + res, verr := ValidateMP4Media(ctx, seg) + require.NoError(t, verr, "segment %d validates (video+audio both present)", i) + dur := time.Duration(res.MediaData.Duration) + require.Greater(t, dur, 200*time.Millisecond, "segment %d duration sane", i) + require.Less(t, dur, 6*time.Second, "segment %d duration not stretched", i) + require.False(t, res.MediaData.Video[0].BFrames, "VideoToolbox capture has no B-frames") + for _, off := range videoPTSDTSOffsets(t, ctx, seg) { + require.Equal(t, time.Duration(0), off, "segment %d: no-reorder stream must keep PTS == DTS — a nonzero offset means ingest invented a reorder delay", i) + } + } + t.Logf("real Mist capture: %d segments, all validated with PTS == DTS", len(segs)) +} + +// makeBFrameAACFMP4 synthesizes an H264 stream WITH B-frames (PTS ≠ DTS) +// alongside AAC audio in a fragmented MP4 — the shape a hardware encoder or +// non-zerolatency x264 push produces, as delivered by MistServer's live .mp4 +// output. Unlike Matroska, MP4 track fragments carry real decode timestamps, +// so the ingest pipeline needs no DTS reconstruction for the fMP4 muxer to mux +// the reordered stream correctly. b-adapt=false forces x264 to actually emit +// the configured B-frames rather than deciding per-scene. +func makeBFrameAACFMP4(t *testing.T, ctx context.Context, seconds int) []byte { + t.Helper() + desc := strings.Join([]string{ + fmt.Sprintf("videotestsrc num-buffers=%d ! video/x-raw,width=320,height=240,framerate=30/1 ! x264enc bframes=2 b-adapt=false key-int-max=30 speed-preset=veryfast ! h264parse ! mp4mux name=mux fragment-duration=500 ! appsink name=sink", seconds*30), + fmt.Sprintf("audiotestsrc num-buffers=%d samplesperbuffer=1024 ! audio/x-raw,rate=48000,channels=2 ! audioconvert ! fdkaacenc ! aacparse ! mux.", seconds*47), + }, "\n") + return runSynthPipeline(t, ctx, desc) +} + +// TestMP4IngestBFramesValidate is the regression test for B-frame ingest: +// every emitted segment must survive the full ValidateMP4Media chokepoint with +// a sane duration, and the stream's real reorder offsets must be preserved. +// (On the old MKV path this scenario originally lost DTS entirely — mp4mux +// stretched the video track ~2.2× and every segment failed validation with +// "no audio in segment"; the h264timestamper fix for that in turn minted +// spurious offsets on no-reorder streams. fMP4's container DTS sidesteps the +// whole trade-off, and this test pins the B-frame half of it.) +func TestMP4IngestBFramesValidate(t *testing.T) { + ctx := context.Background() + mp4 := makeBFrameAACFMP4(t, ctx, 8) + + segs, err := runMP4ThroughIngestWorkerSegments(t, mp4, false) + require.NoError(t, err, "B-frame stream ingests cleanly") + require.GreaterOrEqual(t, len(segs), 6, "B-frame stream emits its segments") + + sawBFrames := false + sawReorderOffset := false + for i, seg := range segs { + res, verr := ValidateMP4Media(ctx, seg) + require.NoError(t, verr, "segment %d validates (video+audio both present)", i) + dur := time.Duration(res.MediaData.Duration) + require.Greater(t, dur, 500*time.Millisecond, "segment %d duration sane", i) + require.Less(t, dur, 2*time.Second, "segment %d duration not stretched", i) + if res.MediaData.Video[0].BFrames { + sawBFrames = true + } + for _, off := range videoPTSDTSOffsets(t, ctx, seg) { + if off > 0 { + sawReorderOffset = true + } + } + } + require.True(t, sawBFrames, "synthesized stream actually contains B-frames — if this fails the test no longer exercises the reorder path") + require.True(t, sawReorderOffset, "B-frame stream keeps its real PTS−DTS reorder offsets through ingest") + t.Logf("B-frame stream: %d segments, all validated", len(segs)) +} diff --git a/pkg/media/mist_pull.go b/pkg/media/mist_pull.go new file mode 100644 index 00000000..528dcad4 --- /dev/null +++ b/pkg/media/mist_pull.go @@ -0,0 +1,154 @@ +package media + +import ( + "bufio" + "context" + "fmt" + "io" + "net" + "net/http" + "net/url" + "strings" + "time" + + "stream.place/streamplace/pkg/log" +) + +// mistPullConnectGrace bounds how long we retry the initial GET while Mist +// boots the freshly-pushed stream: PUSH_REWRITE fires BEFORE Mist accepts the +// push, so the stream may 404 (or refuse the connection) for a moment before +// media flows. Var (not const) so tests can shorten it. +var mistPullConnectGrace = 30 * time.Second + +// mistPullRetryBackoff paces those initial connect retries. +var mistPullRetryBackoff = 500 * time.Millisecond + +// MistPullIngest ingests a live stream by PULLING MistServer's fragmented-MP4 +// HTTP output for mistStreamName — the replacement for the old MKVExec push +// bridge (Mist exec'ing `streamplace live` and POSTing MKV to /live). Pulling +// .mp4 instead of receiving .mkv matters: MP4 fragments carry real decode +// timestamps, so ingest no longer has to reconstruct DTS from the H264 +// bitstream (see buildMP4IngestPipeline). +// +// It's kicked off from the PUSH_REWRITE trigger — the moment we've authed an +// incoming Mist push and minted its signer — and runs for the life of the +// stream: the GET body ends when the push ends. The isolated path hands the +// raw response connection to a detached worker (fd-passing, exactly like the +// old hijacked-POST path), so a main restart doesn't interrupt the ingest. +func (mm *MediaManager) MistPullIngest(ctx context.Context, mistStreamName string, ms MediaSigner) error { + hostport := fmt.Sprintf("127.0.0.1:%d", mm.cli.MistHTTPPort) + // PathEscape leaves '+' (a legal path character) alone, but HTTP servers + // commonly decode it as a space — Mist wildcard names are full of them + // (stream+_), so escape it explicitly. + path := "/" + strings.ReplaceAll(url.PathEscape(mistStreamName), "+", "%2B") + ".mp4" + ctx = log.WithLogValues(ctx, "streamer", ms.Streamer(), "mist-stream", mistStreamName) + + conn, prebuf, chunked, err := mistPullConnect(ctx, hostport, path) + if err != nil { + return fmt.Errorf("mist pull: %w", err) + } + log.Log(ctx, "mist pull connected", "url", hostport+path) + + if mm.cli.IsolatedIngest { + // Zero-downtime path: the detached worker owns the pull connection (so + // it survives a main restart) and serves signed segments back over its + // socket — the same machinery as the old hijacked inbound push, with + // the connection pointing the other way. + return mm.MP4IngestDetached(ctx, conn, prebuf, chunked, ms) + } + defer conn.Close() + body := WorkerInput(IngestWorkerConfig{Prebuf: prebuf, Chunked: chunked}, conn) + return mm.MP4Ingest(ctx, body, ms) +} + +// mistPullConnect dials Mist's HTTP output and issues the GET by hand — not +// through http.Client — because the isolated path needs the raw *net.TCPConn +// to fd-pass to the worker. It consumes the response headers and returns the +// connection positioned at the body, plus any body bytes the header read +// buffered past the headers (prebuf) and whether the body is chunked — the +// same (conn, prebuf, chunked) shape the old hijacked-POST ingest produced, so +// the downstream machinery is shared unchanged. +// +// It retries while Mist boots the stream (mistPullConnectGrace): a refused +// connection or a non-200 just means the push hasn't started flowing yet. +func mistPullConnect(ctx context.Context, hostport, path string) (*net.TCPConn, []byte, bool, error) { + giveUp := time.Now().Add(mistPullConnectGrace) + for { + conn, prebuf, chunked, err := tryMistGET(ctx, hostport, path) + if err == nil { + return conn, prebuf, chunked, nil + } + if time.Now().After(giveUp) { + return nil, nil, false, fmt.Errorf("stream never came up at %s%s: %w", hostport, path, err) + } + select { + case <-ctx.Done(): + return nil, nil, false, ctx.Err() + case <-time.After(mistPullRetryBackoff): + } + } +} + +// tryMistGET is one attempt: dial, send the GET, read the response headers. +// On a 200 it hands back the connection + buffered body bytes; anything else +// is an error and the connection is closed. +func tryMistGET(ctx context.Context, hostport, path string) (*net.TCPConn, []byte, bool, error) { + d := net.Dialer{Timeout: 5 * time.Second} + raw, err := d.DialContext(ctx, "tcp", hostport) + if err != nil { + return nil, nil, false, err + } + conn, ok := raw.(*net.TCPConn) + if !ok { + raw.Close() + return nil, nil, false, fmt.Errorf("expected TCP connection, got %T", raw) + } + // Connection: close — one stream per connection, body runs to EOF (or + // chunked-EOS) when the Mist stream ends. No keepalive reuse to reason about. + req := "GET " + path + " HTTP/1.1\r\n" + + "Host: " + hostport + "\r\n" + + "User-Agent: streamplace-ingest\r\n" + + "Accept: video/mp4\r\n" + + "Connection: close\r\n\r\n" + if err := conn.SetDeadline(time.Now().Add(10 * time.Second)); err != nil { + conn.Close() + return nil, nil, false, err + } + if _, err := io.WriteString(conn, req); err != nil { + conn.Close() + return nil, nil, false, err + } + br := bufio.NewReader(conn) + resp, err := http.ReadResponse(br, nil) + if err != nil { + conn.Close() + return nil, nil, false, err + } + if resp.StatusCode != http.StatusOK { + conn.Close() + return nil, nil, false, fmt.Errorf("mist returned %s", resp.Status) + } + if err := conn.SetDeadline(time.Time{}); err != nil { // clear; streaming has no deadline + conn.Close() + return nil, nil, false, err + } + chunked := false + for _, te := range resp.TransferEncoding { + if te == "chunked" { + chunked = true + } + } + // The header read buffered some raw body bytes; peel them off the bufio so + // the caller can prepend them to the (otherwise unbuffered) connection. + // resp.Body is deliberately never read — it would de-chunk, and the worker + // wants the raw stream + the chunked flag. + var prebuf []byte + if n := br.Buffered(); n > 0 { + prebuf = make([]byte, n) + if _, err := io.ReadFull(br, prebuf); err != nil { + conn.Close() + return nil, nil, false, err + } + } + return conn, prebuf, chunked, nil +} diff --git a/pkg/media/mist_pull_test.go b/pkg/media/mist_pull_test.go new file mode 100644 index 00000000..8316b816 --- /dev/null +++ b/pkg/media/mist_pull_test.go @@ -0,0 +1,113 @@ +package media + +import ( + "context" + "fmt" + "io" + "net" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +// TestMistPullConnect drives mistPullConnect against a fake Mist: the first +// request 404s (the push hasn't started flowing yet — PUSH_REWRITE fires +// before Mist accepts the push), the second streams a chunked body. The +// connector must retry through the 404, then hand back the raw connection + +// buffered bytes + chunked flag such that WorkerInput reconstructs exactly the +// media bytes — the same contract the old hijacked-POST path provided. +func TestMistPullConnect(t *testing.T) { + payload := make([]byte, 256*1024) // big enough to outsize any header-read buffering + for i := range payload { + payload[i] = byte(i) + } + + var calls atomic.Int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + require.Equal(t, "/stream%2Bdid:test:abc_123.mp4", r.URL.EscapedPath()) + if calls.Add(1) == 1 { + http.Error(w, "stream not ready", http.StatusNotFound) + return + } + w.Header().Set("Content-Type", "video/mp4") + fl := w.(http.Flusher) + // Stream in pieces with flushes so Go's server chunks the response — + // the shape Mist's live output has. + for i := 0; i < len(payload); i += 4096 { + end := min(i+4096, len(payload)) + _, err := w.Write(payload[i:end]) + require.NoError(t, err) + fl.Flush() + } + })) + defer srv.Close() + + oldGrace, oldBackoff := mistPullConnectGrace, mistPullRetryBackoff + mistPullConnectGrace, mistPullRetryBackoff = 5*time.Second, 10*time.Millisecond + defer func() { mistPullConnectGrace, mistPullRetryBackoff = oldGrace, oldBackoff }() + + hostport := srv.Listener.Addr().String() + conn, prebuf, chunked, err := mistPullConnect(context.Background(), hostport, "/stream%2Bdid:test:abc_123.mp4") + require.NoError(t, err) + defer conn.Close() + require.GreaterOrEqual(t, calls.Load(), int32(2), "connector retried through the 404") + require.True(t, chunked, "streamed live body is chunked") + + got, err := io.ReadAll(WorkerInput(IngestWorkerConfig{Prebuf: prebuf, Chunked: chunked}, conn)) + require.NoError(t, err) + require.Equal(t, payload, got, "prebuf + raw conn de-frames to the exact media bytes") +} + +// TestMistPullConnectGivesUp: a stream that never comes up must fail within +// the connect grace instead of retrying forever. +func TestMistPullConnectGivesUp(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "no such stream", http.StatusNotFound) + })) + defer srv.Close() + + oldGrace, oldBackoff := mistPullConnectGrace, mistPullRetryBackoff + mistPullConnectGrace, mistPullRetryBackoff = 300*time.Millisecond, 20*time.Millisecond + defer func() { mistPullConnectGrace, mistPullRetryBackoff = oldGrace, oldBackoff }() + + _, _, _, err := mistPullConnect(context.Background(), srv.Listener.Addr().String(), "/nope.mp4") + require.Error(t, err) + require.Contains(t, err.Error(), "never came up") +} + +// TestMistPullConnectRefusedThenUp: Mist itself may not even be listening yet +// (or between restarts); a refused connection is retried like a 404. +func TestMistPullConnectRefusedThenUp(t *testing.T) { + // Reserve a port, then close the listener so the first dials are refused. + l, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + hostport := l.Addr().String() + require.NoError(t, l.Close()) + + oldGrace, oldBackoff := mistPullConnectGrace, mistPullRetryBackoff + mistPullConnectGrace, mistPullRetryBackoff = 5*time.Second, 20*time.Millisecond + defer func() { mistPullConnectGrace, mistPullRetryBackoff = oldGrace, oldBackoff }() + + go func() { + time.Sleep(200 * time.Millisecond) + l2, lerr := net.Listen("tcp", hostport) + if lerr != nil { + return // port raced away; the test will fail on the connect error + } + _ = http.Serve(l2, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + fmt.Fprint(w, "media") + w.(http.Flusher).Flush() + })) + }() + + conn, prebuf, chunked, err := mistPullConnect(context.Background(), hostport, "/late.mp4") + require.NoError(t, err) + defer conn.Close() + got, err := io.ReadAll(WorkerInput(IngestWorkerConfig{Prebuf: prebuf, Chunked: chunked}, conn)) + require.NoError(t, err) + require.Equal(t, []byte("media"), got) +} diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mp4_ingest.go similarity index 60% rename from pkg/media/mkv_ingest.go rename to pkg/media/mp4_ingest.go index 64182fcf..77314fda 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mp4_ingest.go @@ -14,24 +14,25 @@ import ( "stream.place/streamplace/pkg/log" ) -// ingest a H264+AAC MKV stream (prolly from an RTMP server) -func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms MediaSigner) error { +// ingest a H264+AAC fragmented-MP4 stream (the MistServer live .mp4 output, or +// an fMP4 push to /live) +func (mm *MediaManager) MP4Ingest(ctx context.Context, input io.Reader, ms MediaSigner) error { shouldRecord, err := mm.shouldRecord(ctx, ms.Streamer()) if err != nil { return err } if shouldRecord { - log.Log(ctx, "recording RTMP stream to file", "streamer", ms.Streamer()) + log.Log(ctx, "recording ingest stream to file", "streamer", ms.Streamer()) pr, pw := io.Pipe() input = io.TeeReader(input, pw) go func() { - err := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mkv") + err := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mp4") if err != nil { log.Error(ctx, "error dumping to file", "error", err) } }() } else { - log.Log(ctx, "not recording RTMP stream to file", "streamer", ms.Streamer()) + log.Log(ctx, "not recording ingest stream to file", "streamer", ms.Streamer()) } ctx, cancel := context.WithCancel(ctx) defer cancel() @@ -40,7 +41,7 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media if err != nil { return err } - pipeline, err := buildMKVIngestPipeline(ctx, input, signer) + pipeline, err := buildMP4IngestPipeline(ctx, input, signer) if err != nil { return err } @@ -64,41 +65,42 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media return <-busErr } -// buildMKVIngestPipeline builds the H264+AAC MKV demux graph (video → h264parse, -// audio → Opus re-encode) and links both branches into signerElem — the muxl -// signing bin that emits one bare canonical .m4s per GoP. Shared by the -// in-process MKVIngest and the isolated ingest worker, which differ only in +// buildMP4IngestPipeline builds the H264+AAC fragmented-MP4 demux graph (video +// → h264parse, audio → Opus re-encode) and links both branches into signerElem +// — the muxl signing bin that emits one bare canonical .m4s per GoP. Shared by +// the in-process MP4Ingest and the isolated ingest worker, which differ only in // where signerElem routes its segments (ValidateMP4 vs. a frame writer to the // main process). -func buildMKVIngestPipeline(ctx context.Context, input io.Reader, signerElem *gst.Element) (*gst.Pipeline, error) { - // Queue sizing: matroskademux feeds both branches from one thread, and the - // fMP4 muxer downstream is an aggregator — it consumes NOTHING until every - // pad has data. If the video track goes sparse (e.g. MistServer drops all - // delta frames when a push falls behind, leaving ~1s-apart keyframes), the - // audio branch must buffer a full video-frame gap while the demux walks the - // byte stream to the next video frame. gst's default queue caps at +// +// The source is fMP4, not MKV, very much on purpose: MP4 track fragments carry +// both decode (tfdt/trun) and presentation (ctts) timestamps, so qtdemux hands +// us the encoder's real DTS. Matroska carries only presentation timestamps, so +// the old MKV ingest had to *reconstruct* DTS with h264timestamper — which +// guesses a worst-case full-DPB reorder window for streams whose SPS doesn't +// declare one (notably VideoToolbox), minting a constant spurious PTS−DTS +// offset that pushed every GoP's presentation past its segment's declared +// window and broke WebRTC playback at every keyframe. Real DTS in the +// container means no reconstruction and no guessing. +func buildMP4IngestPipeline(ctx context.Context, input io.Reader, signerElem *gst.Element) (*gst.Pipeline, error) { + // Queue sizing: qtdemux feeds both branches from one thread, and the fMP4 + // muxer downstream is an aggregator — it consumes NOTHING until every pad + // has data. If the video track goes sparse (e.g. MistServer drops all delta + // frames when a push falls behind, leaving ~1s-apart keyframes), the audio + // branch must buffer a full video-frame gap while the demux walks the byte + // stream to the next video frame. gst's default queue caps at // max-size-time=1s, so a ≥1s video gap fills the audio queue, blocks the // demux, starves the muxer's video pad, and deadlocks the whole graph with // no EOS — a live stream wedges until the watchdog kills it. Use the shared // Queue2Big preset (no time/buffer cap, generous byte cap) like the other // demux-fed pipelines (transcode, rtmp_push, packetize, media_data_parser). - // - // h264timestamper: Matroska blocks carry only presentation timestamps, so - // for B-frame streams (PTS ≠ DTS) matroskademux emits reordered PTS with - // dts=none — and h264parse does not reconstruct DTS. The fMP4 muxer needs - // DTS to mux a reordered stream; without it it treats the jumbled PTS as - // monotonic timing and stretches the video track (~2.2× on a real capture), - // which downstream makes qtdemux EOS the audio pad early and every segment - // fails validation with "no audio in segment". h264timestamper rebuilds - // DTS from the H264 picture order count. pipelineSlice := []string{ - "appsrc name=streamsrc ! matroskademux name=demux", - "demux. ! " + constants.Queue2Big + " ! h264parse ! h264timestamper name=videoout", + "appsrc name=streamsrc ! qtdemux name=demux", + "demux. ! " + constants.Queue2Big + " ! h264parse name=videoout", "demux. ! " + constants.Queue2Big + " ! fdkaacdec ! audioresample ! opusenc name=audioenc", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { - return nil, fmt.Errorf("error creating MKVIngest pipeline: %w", err) + return nil, fmt.Errorf("error creating MP4Ingest pipeline: %w", err) } srcele, err := pipeline.GetElementByName("streamsrc") if err != nil { @@ -132,7 +134,7 @@ func (mm *MediaManager) dumpToFile(ctx context.Context, r io.Reader, user string filename := fmt.Sprintf("%s%s", now.FileSafeString(), filesuffix) // Streams to S3 when configured (production), else a local file under DataDir // (dev). Close finalizes either target — for S3 it commits the upload. - f, err := mm.cli.DebugRecordingCreate(ctx, []string{"debug-recordings", user, filename}, "video/x-matroska", false) + f, err := mm.cli.DebugRecordingCreate(ctx, []string{"debug-recordings", user, filename}, "video/mp4", false) if err != nil { return fmt.Errorf("failed to create debug recording: %w", err) } diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index 65fdabbc..9269efb2 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -236,7 +236,7 @@ func SegmentUnsigned(ctx context.Context, cli *config.CLI, streamer string, inpu } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) if err != nil { - return fmt.Errorf("error creating MKVIngest pipeline: %w", err) + return fmt.Errorf("error creating SegmentUnsigned pipeline: %w", err) } srcele, err := pipeline.GetElementByName("appsrc") diff --git a/pkg/media/whip_worker.go b/pkg/media/whip_worker.go index de598bcb..be8ad040 100644 --- a/pkg/media/whip_worker.go +++ b/pkg/media/whip_worker.go @@ -14,7 +14,7 @@ import ( ) // ServeWHIPIngestWorkerSocket is the WHIP counterpart of -// ServeMKVIngestWorkerSocket. Unlike MKV there's no socket/fd to pass in: the +// ServeMP4IngestWorkerSocket. Unlike the fMP4 worker there is no socket/fd to pass in: the // worker OWNS the PeerConnection, so it creates it from cfg.OfferSDP (binding its // own UDP sockets), generates the SDP answer, and emits it as the FIRST frame on // the unix socket — main reads that Answer frame and returns it to the WHIP diff --git a/pkg/media/whip_worker_test.go b/pkg/media/whip_worker_test.go index 717ebe96..481e3496 100644 --- a/pkg/media/whip_worker_test.go +++ b/pkg/media/whip_worker_test.go @@ -47,7 +47,7 @@ func whipClientOffer(t *testing.T) (*webrtc.PeerConnection, *webrtc.TrackLocalSt // an offer it builds the PeerConnection, generates the SDP answer, and emits it // as the FIRST frame on its socket — the synchronous reply main returns to the // WHIP client. (Media flow → signed segments rides the same webRTCIngestPipeline -// the in-process path uses, plus the transcoder/frame machinery the MKV tests +// the in-process path uses, plus the transcoder/frame machinery the fMP4 ingest tests // already cover.) func TestWHIPWorkerAnswersOffer(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) @@ -149,7 +149,7 @@ func produceWHIPMedia(t *testing.T, ctx context.Context, video, audio *webrtc.Tr // TestWHIPWorkerLoopback is the full WHIP media path: a pion client offers, // connects to the worker (which owns the PeerConnection), and streams real // H264+Opus RTP; the worker must mux+sign+transcode it and serve a valid signed -// dual-codec segment over its socket — the WHIP parity of the MKV worker e2e +// dual-codec segment over its socket — the WHIP parity of the fMP4 ingest worker e2e // test. func TestWHIPWorkerLoopback(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) diff --git a/pkg/media/worker_watchdog.go b/pkg/media/worker_watchdog.go index ea5cd176..8d3ec3ff 100644 --- a/pkg/media/worker_watchdog.go +++ b/pkg/media/worker_watchdog.go @@ -15,7 +15,7 @@ import ( // 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, +// This is the worker-side counterpart to MP4IngestIsolated'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. -- 2.51.2 From 857036b7cd7c7458773fd3e0cbff376844356287 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 22 Jul 2026 13:06:21 -0700 Subject: [PATCH 2/3] media: serialize the mist pull request with net/http, not by hand MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tryMistGET still dials its own conn — that part is load-bearing (the raw fd gets passed to the detached worker, which http.Client can't provide) — but the request itself is now a real *http.Request written with req.Write instead of a concatenated header string, and http.ReadResponse gets the request for context. Same wire bytes, stdlib framing. The %2B-escaped wildcard path survives req.Write via URL.RawPath (covered by TestMistPullConnect's EscapedPath assertion). Co-Authored-By: Claude Fable 5 --- pkg/media/mist_pull.go | 21 +++++++++++++-------- 1 file changed, 13 insertions(+), 8 deletions(-) diff --git a/pkg/media/mist_pull.go b/pkg/media/mist_pull.go index 528dcad4..66d473be 100644 --- a/pkg/media/mist_pull.go +++ b/pkg/media/mist_pull.go @@ -103,23 +103,28 @@ func tryMistGET(ctx context.Context, hostport, path string) (*net.TCPConn, []byt raw.Close() return nil, nil, false, fmt.Errorf("expected TCP connection, got %T", raw) } - // Connection: close — one stream per connection, body runs to EOF (or + // A real *http.Request serialized by the stdlib — we only own the conn by + // hand (it gets fd-passed to the worker), not the HTTP framing. req.Close + // sends Connection: close: one stream per connection, body runs to EOF (or // chunked-EOS) when the Mist stream ends. No keepalive reuse to reason about. - req := "GET " + path + " HTTP/1.1\r\n" + - "Host: " + hostport + "\r\n" + - "User-Agent: streamplace-ingest\r\n" + - "Accept: video/mp4\r\n" + - "Connection: close\r\n\r\n" + req, err := http.NewRequestWithContext(ctx, http.MethodGet, "http://"+hostport+path, nil) + if err != nil { + conn.Close() + return nil, nil, false, err + } + req.Close = true + req.Header.Set("User-Agent", "streamplace-ingest") + req.Header.Set("Accept", "video/mp4") if err := conn.SetDeadline(time.Now().Add(10 * time.Second)); err != nil { conn.Close() return nil, nil, false, err } - if _, err := io.WriteString(conn, req); err != nil { + if err := req.Write(conn); err != nil { conn.Close() return nil, nil, false, err } br := bufio.NewReader(conn) - resp, err := http.ReadResponse(br, nil) + resp, err := http.ReadResponse(br, req) if err != nil { conn.Close() return nil, nil, false, err -- 2.51.2 From aa11bba08fe27294cc0dda308ebdd9b33aaf94d0 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 22 Jul 2026 13:20:10 -0700 Subject: [PATCH 3/3] media: diagnose legacy MKV pushes; default mist-http-port to 28080 Field report from the first live test of the fMP4 ingest: a MistServer container still running the legacy MKVExec process config POSTed MKV to /live on a 4s restart loop, and each attempt died as an instant cryptic qtdemux failure (a pile of ~200-byte truncated debug recordings). Meanwhile the pull ingest dialed the old default port 18080 while Mist listened on 28080, so it never connected at all. - buildMP4IngestPipeline now peeks the stream and rejects the EBML magic with an error that names the actual problem (legacy MKVExec config) instead of letting qtdemux die confusingly. - mist-http-port default 18080 -> 28080, matching docker/mistserver.json (the generated dev config derives Mist's listener from the same flag, so both worlds stay consistent). Co-Authored-By: Claude Fable 5 --- pkg/config/config.go | 4 ++-- pkg/media/mist_mp4_ingest_test.go | 12 ++++++++++++ pkg/media/mp4_ingest.go | 32 +++++++++++++++++++++++++++++++ 3 files changed, 46 insertions(+), 2 deletions(-) diff --git a/pkg/config/config.go b/pkg/config/config.go index 1ba3931b..03a78af3 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -1032,8 +1032,8 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }) cmd.Flags = append(cmd.Flags, &urfavecli.IntFlag{ Name: "mist-http-port", - Usage: "MistServer HTTP port (internal use only)", - Value: 18080, + Usage: "MistServer HTTP port (internal use only) — ingest pulls Mist's live fMP4 output from this port, so it must match the running Mist config (docker/mistserver.json uses 28080, the default here)", + Value: 28080, Destination: &cli.MistHTTPPort, Sources: urfavecli.EnvVars("SP_MIST_HTTP_PORT"), }) diff --git a/pkg/media/mist_mp4_ingest_test.go b/pkg/media/mist_mp4_ingest_test.go index 13a60bcc..d933ca18 100644 --- a/pkg/media/mist_mp4_ingest_test.go +++ b/pkg/media/mist_mp4_ingest_test.go @@ -282,3 +282,15 @@ func TestMP4IngestBFramesValidate(t *testing.T) { require.True(t, sawReorderOffset, "B-frame stream keeps its real PTS−DTS reorder offsets through ingest") t.Logf("B-frame stream: %d segments, all validated", len(segs)) } + +// TestMP4IngestRejectsMatroskaWithDiagnosis: an MKV stream landing on the +// fMP4 ingest (a MistServer still running the legacy MKVExec process config) +// must fail immediately with a message that names the actual problem — not a +// generic qtdemux parse error on an endless Mist-side restart loop. +func TestMP4IngestRejectsMatroskaWithDiagnosis(t *testing.T) { + mkvish := append(append([]byte{}, matroskaMagic...), make([]byte, 1024)...) + _, err := runMP4ThroughIngestWorker(t, mkvish, false) + require.Error(t, err) + require.Contains(t, err.Error(), "Matroska") + require.Contains(t, err.Error(), "MKVExec") +} diff --git a/pkg/media/mp4_ingest.go b/pkg/media/mp4_ingest.go index 77314fda..776e804b 100644 --- a/pkg/media/mp4_ingest.go +++ b/pkg/media/mp4_ingest.go @@ -1,7 +1,10 @@ package media import ( + "bufio" + "bytes" "context" + "errors" "fmt" "io" "strings" @@ -81,7 +84,36 @@ func (mm *MediaManager) MP4Ingest(ctx context.Context, input io.Reader, ms Media // offset that pushed every GoP's presentation past its segment's declared // window and broke WebRTC playback at every keyframe. Real DTS in the // container means no reconstruction and no guessing. +// matroskaMagic is the EBML header every Matroska/WebM stream opens with. +var matroskaMagic = []byte{0x1A, 0x45, 0xDF, 0xA3} + +// rejectMatroska peeks at the ingest stream and fails fast with a diagnosis if +// it's Matroska. MKV was this pipeline's previous bridge format, so the most +// likely stray MKV source is a MistServer still running the legacy MKVExec +// process config (`streamplace live` POSTing MKV to /live on a restart loop) — +// without the sniff that just looks like qtdemux dying instantly, over and +// over, which is a miserable thing to debug. Returns a reader that includes +// the peeked bytes. +func rejectMatroska(input io.Reader) (io.Reader, error) { + br := bufio.NewReader(input) + head, err := br.Peek(len(matroskaMagic)) + if err != nil { + if errors.Is(err, io.EOF) { + return br, nil // shorter than the magic; let the pipeline EOS/complain + } + return nil, fmt.Errorf("peek ingest stream: %w", err) + } + if bytes.Equal(head, matroskaMagic) { + return nil, fmt.Errorf("ingest input is Matroska (MKV), but this node ingests fragmented MP4 — a MistServer running the legacy MKVExec process config is probably still pushing MKV to /live; update its config (see docker/mistserver.json)") + } + return br, nil +} + func buildMP4IngestPipeline(ctx context.Context, input io.Reader, signerElem *gst.Element) (*gst.Pipeline, error) { + input, err := rejectMatroska(input) + if err != nil { + return nil, err + } // Queue sizing: qtdemux feeds both branches from one thread, and the fMP4 // muxer downstream is an aggregator — it consumes NOTHING until every pad // has data. If the video track goes sparse (e.g. MistServer drops all delta