diff --git a/pkg/media/transcode.go b/pkg/media/transcode.go index 0dfa2d962..09746d664 100644 --- a/pkg/media/transcode.go +++ b/pkg/media/transcode.go @@ -10,6 +10,7 @@ import ( "github.com/go-gst/go-gst/gst" "stream.place/streamplace/pkg/atproto" + "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/crypto/signers" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/muxl" @@ -113,19 +114,28 @@ func (mm *MediaManager) audioCompletionTarget(ctx context.Context, seg []byte) ( // track. Callers wire the `src`/`sink` callbacks and set the pipeline playing. // target is the codec being produced: "opus" (source AAC) or "aac" (source Opus). func buildAudioTranscodePipeline(target string) (*gst.Pipeline, error) { + // Queue sizing: a single qtdemux feeds both branches, so if either queue + // hits a limit and blocks the demux, the sibling branch starves and mp4mux + // deadlocks waiting for it (then appsrc backpressures and the Feed blocks). + // gst's default queue caps at max-size-time=1s, which a long GoP overflows: + // a 2s GoP overflowed the video queue's time cap and wedged a live stream + // (TestStreamTranscoderDoubleGopWedge). Use the shared Queue2Big preset + // (no time/buffer cap, generous byte cap) like the other demux-fed pipelines + // (rtmp_push, packetize, media_data_parser) so an over-long GoP flows + // through instead of deadlocking, while memory stays bounded. var audioChain string switch target { case "opus": // source is AAC - audioChain = "queue name=aq ! aacparse ! fdkaacdec ! audioconvert ! audioresample ! opusenc name=aenc" + audioChain = constants.Queue2Big + " name=aq ! aacparse ! fdkaacdec ! audioconvert ! audioresample ! opusenc name=aenc" case "aac": // source is Opus - audioChain = "queue name=aq ! opusparse ! opusdec ! audioconvert ! audioresample ! fdkaacenc name=aenc" + audioChain = constants.Queue2Big + " name=aq ! opusparse ! opusdec ! audioconvert ! audioresample ! fdkaacenc name=aenc" default: return nil, fmt.Errorf("unsupported transcode target %q", target) } pipeline, err := gst.NewPipelineFromString(strings.Join([]string{ "appsrc name=src ! qtdemux name=demux", - "queue name=vq ! h264parse name=vparse", + constants.Queue2Big + " name=vq ! h264parse name=vparse", audioChain, }, "\n")) if err != nil { diff --git a/pkg/media/transcode_stream.go b/pkg/media/transcode_stream.go index 25bb2c808..d920a3635 100644 --- a/pkg/media/transcode_stream.go +++ b/pkg/media/transcode_stream.go @@ -54,6 +54,7 @@ type streamTranscoder struct { feedMu sync.Mutex // serializes Feed so segments enter in order reaper *time.Timer // idle reaper; reset on Feed (set by the registry) + fedSeq int // per-instance feed counter (debug dump ordering); under feedMu mu sync.Mutex started bool @@ -170,6 +171,17 @@ func (t *streamTranscoder) Feed(src []byte, token any) error { t.feedMu.Lock() defer t.feedMu.Unlock() + // Debug: capture each source segment fed to THIS transcoder instance, in + // feed order, so a wedging input sequence can be replayed into a + // streamTranscoder in a test. No-op unless --segment-debug-dir is set. The + // per-instance sequence number is baked into the name (not just the dump's + // async wall-clock stamp) so a burst of sub-second feeds still sorts in feed + // order. Numbering is per-instance, so the run that actually wedges (before + // a reset rebuilds a fresh transcoder) is a self-contained 0..N sequence. + seq := t.fedSeq + t.fedSeq++ + t.mm.cli.DumpDebugSegment(t.ctx, fmt.Sprintf("transcoder-feed-%s-%05d.m4s", t.target, seq), bytes.NewReader(src)) + t.mu.Lock() switch { case t.closed: diff --git a/pkg/media/transcode_wedge_test.go b/pkg/media/transcode_wedge_test.go new file mode 100644 index 000000000..bd2b4d0c8 --- /dev/null +++ b/pkg/media/transcode_wedge_test.go @@ -0,0 +1,86 @@ +package media + +import ( + "context" + "fmt" + "os" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/crypto/signers" + "stream.place/streamplace/test/remote" +) + +// TestStreamTranscoderDoubleGopWedge replays the captured input that wedged a +// live stream: three normal ~1s GoPs to prime the continuous encoder, then one +// double-length (~2s, 120 video / 100 audio samples) GoP. That double-length +// GoP was the sole ~2s inter-arrival in a 10-minute capture, and it's exactly +// where ingest froze. +// +// The failure mode it guards against: gst's default queue caps at +// max-size-time=1s, so a 2s GoP overflowed the video queue's time cap, blocked +// qtdemux (which feeds both branches), starved the audio branch, and +// deadlocked mp4mux — backpressuring appsrc until the synchronous Feed +// (feedW.Write into an unbuffered io.Pipe) blocked, wedging the serial ingest +// path. Before the Queue2Big sizing fix this Fataled on the 25s timeout +// (completing only 2/4); with it the 2s GoP flows through. +func TestStreamTranscoderDoubleGopWedge(t *testing.T) { + // Captured run-2 window: three 1s primes then the 2s double GoP (#185). + files := []string{ + remote.RemoteFixture("1c0c04ea0f96a6abbaaf9985ea3691c3ac54728f0a52e0575eb5968e61ca30b2/2026-05-28T18-34-19-607Z-transcoder-feed-aac-00182.m4s"), + remote.RemoteFixture("536cf3325314e5ecbd88ac482ec5253ee781cc47298de64b2a1820e156d7ba83/2026-05-28T18-34-20-616Z-transcoder-feed-aac-00183.m4s"), + remote.RemoteFixture("7b3e0d0e33e61eadb549aea9a60347714ed6182858c1e39a2033e4e00383c881/2026-05-28T18-34-21-610Z-transcoder-feed-aac-00184.m4s"), + remote.RemoteFixture("caada3d7c3226a9fb7ed5b9b20e34efbf7154d8662254a8cbbf33ecff5ca497d/2026-05-28T18-34-23-645Z-transcoder-feed-aac-00185.m4s"), + } + + ctx := context.Background() + ms := newBareSegmentSigner(t) + mm := &MediaManager{cli: &config.CLI{BroadcasterHost: "test.example.com"}} + keyPEM, err := signers.MarshalES256KPrivateKeyPEM(ms.Signer) + require.NoError(t, err) + + var mu sync.Mutex + completed := 0 + tr := mm.newStreamTranscoder(ctx, "aac", ms.Cert, keyPEM, func(_ any, _ []byte) { + mu.Lock() + completed++ + mu.Unlock() + }) + + // Feed in a goroutine so the wedge surfaces as a timeout rather than a hang: + // Feed blocks (feedW.Write) when the pipeline deadlocks. + done := make(chan error, 1) + go func() { + for i, f := range files { + b, rerr := os.ReadFile(f) + if rerr != nil { + done <- rerr + return + } + if ferr := tr.Feed(b, i); ferr != nil { + done <- fmt.Errorf("feed %d: %w", i, ferr) + return + } + } + done <- tr.Close() + }() + + select { + case err := <-done: + require.NoError(t, err) + case <-time.After(25 * time.Second): + mu.Lock() + n := completed + mu.Unlock() + t.Fatalf("WEDGED: feed+close did not finish in 25s (completed %d/%d) — "+ + "a double-length GoP stalled the transcode pipeline and blocked Feed", n, len(files)) + } + + // All four GoPs complete: the three primes plus the 2s GoP (flushed by Close). + mu.Lock() + defer mu.Unlock() + require.Equal(t, len(files), completed, "every fed GoP should complete (incl. the 2s double GoP)") +}