package media import ( "bytes" "context" "errors" "fmt" "io" "net/http" "net/http/httptest" "os" "path/filepath" "strings" "sync" "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/pkg/s3" ) // TestWorkerInputDeframes checks the body-deframing the worker applies to the // 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-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)) // prebuf = the slice main already read; the rest is still on the fd. cfg := IngestWorkerConfig{Chunked: true, Prebuf: append([]byte(nil), body[:5]...)} got, err := io.ReadAll(WorkerInput(cfg, bytes.NewReader(body[5:]))) require.NoError(t, err) require.Equal(t, payload, got, "prebuf + chunked fd de-frames to the original media") // Raw (no prebuf, not chunked) passes through unchanged. got, err = io.ReadAll(WorkerInput(IngestWorkerConfig{}, bytes.NewReader(payload))) require.NoError(t, err) require.Equal(t, payload, got) } // 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 ! mp4mux name=mux fragment-duration=500 ! appsink name=sink", "d. ! queue ! opusdec ! audioconvert ! audioresample ! fdkaacenc ! aacparse ! mux.", }, "\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, "remux to fragmented H264+AAC MP4") require.NotEmpty(t, buf.Bytes(), "remux produced an fMP4") return buf.Bytes() } // 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, 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 ! mp4mux fragment-duration=500 ! appsink name=sink", seconds*47) 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 audio-only fMP4") require.NotEmpty(t, buf.Bytes()) return buf.Bytes() } // 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 TestRunMP4IngestWorkerProducesValidSignedFrames(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) require.NoError(t, err) manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) require.NoError(t, err) // Provide the node transcode signer (the same test key serves both roles // here, as in the transcoder tests) so the worker completes to dual-codec. cfg := IngestWorkerConfig{ StreamerDID: ms.Streamer(), KeyPEM: keyPEM, CertPEM: ms.Cert, Manifest: manifest, NodeCertPEM: ms.Cert, NodeKeyPEM: keyPEM, BroadcasterHost: "test.example.com", } mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) // 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, RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), frames, func() []byte { return cfg.Manifest })) r := ingestframe.NewReader(&buf) var segs int for { typ, payload, err := r.ReadFrame() if errors.Is(err, io.EOF) { break } require.NoError(t, err) require.Equal(t, ingestframe.Segment, typ, "worker emits only Segment frames; End is the subcommand's job") require.NotEmpty(t, payload) out, err := muxl.RunMuxlVerify(ctx, bytes.NewReader(payload)) require.NoError(t, err, "segment %d verify", segs) require.NotContains(t, out, `"validation_state":"Invalid"`, "segment %d must validate", segs) // With a node key the worker completes to dual-codec: every segment must // carry both the source Opus and a worker-transcoded AAC track. codecs := audioCodecsOf(t, ctx, payload) hasOpus, hasAAC := false, false for _, c := range codecs { if isOpusCodec(c) { hasOpus = true } if isAACCodec(c) { hasAAC = true } } require.True(t, hasOpus, "segment %d keeps source Opus (got %v)", segs, codecs) require.True(t, hasAAC, "segment %d gains worker-transcoded AAC (got %v)", segs, codecs) segs++ } require.GreaterOrEqual(t, segs, 1, "worker emitted at least one signed dual-codec segment") t.Logf("worker emitted %d valid dual-codec segments", segs) } // 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.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 TestRunMP4IngestWorkerRecords(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) require.NoError(t, err) manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) require.NoError(t, err) dataDir := t.TempDir() cfg := IngestWorkerConfig{ StreamerDID: ms.Streamer(), KeyPEM: keyPEM, CertPEM: ms.Cert, Manifest: manifest, BroadcasterHost: "test.example.com", Record: true, DataDir: dataDir, } mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) // Frames are irrelevant here (recording is on the input side); discard them. 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.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.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, mp4) }, 10*time.Second, 25*time.Millisecond, "worker records the ingest media verbatim") } // fakeS3Server is a minimal path-style S3 endpoint speaking just enough of the // multipart-upload protocol for UploadWriter: initiate → upload parts → // complete. Completed objects land in objects keyed by "/". type fakeS3Server struct { mu sync.Mutex parts map[string][]byte // "#" → body objects map[string][]byte // completed "/" → body } func newFakeS3Server() *fakeS3Server { return &fakeS3Server{parts: map[string][]byte{}, objects: map[string][]byte{}} } func (f *fakeS3Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { f.mu.Lock() defer f.mu.Unlock() path := strings.TrimPrefix(r.URL.Path, "/") q := r.URL.Query() switch { case r.Method == "POST" && q.Has("uploads"): fmt.Fprintf(w, `test-upload`) case r.Method == "PUT" && q.Has("partNumber"): body, _ := io.ReadAll(r.Body) f.parts[path+"#"+q.Get("partNumber")] = body w.Header().Set("ETag", `"part-`+q.Get("partNumber")+`"`) case r.Method == "POST" && q.Has("uploadId"): var buf []byte for i := 1; ; i++ { part, ok := f.parts[fmt.Sprintf("%s#%d", path, i)] if !ok { break } buf = append(buf, part...) } f.objects[path] = buf fmt.Fprintf(w, `%s`, path) default: w.WriteHeader(http.StatusBadRequest) } } // objectWithPrefix finds a completed object whose key starts with prefix and // ends with suffix (the recording's timestamped filename isn't predictable). func (f *fakeS3Server) objectWithPrefix(prefix, suffix string) ([]byte, bool) { f.mu.Lock() defer f.mu.Unlock() for path, b := range f.objects { if strings.HasPrefix(path, prefix) && strings.HasSuffix(path, suffix) { return b, true } } return nil, false } // TestRunMP4IngestWorkerRecordsToS3 proves the debug recording streams to S3 // when main hands its S3 config over the handshake (cfg.S3) — the production // shape. Without that plumbing the worker's minimal CLI has no S3 fields and // DebugRecordingCreate silently falls back to local disk, which is exactly the // regression this guards against: recordings must land in the bucket, not under // DataDir. func TestRunMP4IngestWorkerRecordsToS3(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) require.NoError(t, err) manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) require.NoError(t, err) fake := newFakeS3Server() srv := httptest.NewServer(fake) defer srv.Close() dataDir := t.TempDir() cfg := IngestWorkerConfig{ StreamerDID: ms.Streamer(), KeyPEM: keyPEM, CertPEM: ms.Cert, Manifest: manifest, BroadcasterHost: "test.example.com", Record: true, DataDir: dataDir, S3: &s3.Config{ Endpoint: srv.URL, Bucket: "debug-bucket", AccessKeyID: "test-access", SecretAccessKey: "test-secret", Region: "auto", }, } mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) require.NoError(t, RunMP4IngestWorker(ctx, cfg, bytes.NewReader(mp4), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest })) // The upload commits asynchronously (the dump goroutine's Close); the object // must appear at debug-bucket/debug-recordings//.rtmp.mp4 holding // exactly the ingested media. wantPrefix := "debug-bucket/debug-recordings/" + ms.Streamer() + "/" require.Eventually(t, func() bool { got, ok := fake.objectWithPrefix(wantPrefix, ".rtmp.mp4") return ok && bytes.Equal(got, mp4) }, 10*time.Second, 25*time.Millisecond, "worker streams the recording to the S3 bucket verbatim") // And nothing fell back to local disk. matches, _ := filepath.Glob(filepath.Join(dataDir, "debug-recordings", "*", "*")) require.Empty(t, matches, "recording must go to S3, not DataDir") } // 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 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 TestMP4IngestIsolatedWedgeContained.) func TestRunMP4IngestWorkerSelfWatchdog(t *testing.T) { old := ingestWorkerWatchdog ingestWorkerWatchdog = 3 * time.Second defer func() { ingestWorkerWatchdog = old }() ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) require.NoError(t, err) manifest, err := ms.buildManifest(ctx, time.Now().UnixMilli()) require.NoError(t, err) cfg := IngestWorkerConfig{ StreamerDID: ms.Streamer(), KeyPEM: keyPEM, CertPEM: ms.Cert, Manifest: manifest, BroadcasterHost: "test.example.com", } wedge := makeAudioOnlyAACFMP4(t, ctx, 5) start := time.Now() done := make(chan error, 1) go func() { done <- RunMP4IngestWorker(ctx, cfg, bytes.NewReader(wedge), ingestframe.NewWriter(io.Discard), func() []byte { return cfg.Manifest }) }() select { case <-done: elapsed := time.Since(start) require.Less(t, elapsed, 25*time.Second, "watchdog bounded the wedge") t.Logf("worker self-terminated on wedge in %s", elapsed.Round(time.Second)) case <-time.After(30 * time.Second): t.Fatal("worker-side watchdog did not contain the wedge (RunMP4IngestWorker hung)") } } // 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 TestRunMP4IngestWorkerSignsWithSuppliedManifest(t *testing.T) { ctx := context.Background() ms := newBareSegmentSigner(t) keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) require.NoError(t, err) cfg := IngestWorkerConfig{ StreamerDID: ms.Streamer(), KeyPEM: keyPEM, CertPEM: ms.Cert, BroadcasterHost: "test.example.com", } prelive := ms.PrebuiltManifest // c2pa.created only → unpublished live := bytes.Replace(prelive, []byte(`{"action":"c2pa.created"}`), []byte(`{"action":"c2pa.created"},{"action":"c2pa.published"}`), 1) require.NotEqual(t, string(prelive), string(live), "the live manifest must add c2pa.published") mp4 := makeH264AACFMP4(t, ctx, getFixture("5sec.mp4")) firstSegmentPublished := func(manifest []byte) bool { var buf bytes.Buffer 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() require.NoError(t, rerr) require.Equal(t, ingestframe.Segment, typ) res, verr := ValidateMP4Media(ctx, payload) require.NoError(t, verr) return res.Meta.Published } require.False(t, firstSegmentPublished(prelive), "pre-live manifest → unpublished segments") require.True(t, firstSegmentPublished(live), "live manifest → published segments (the fix)") }