From 35e2d82174527843ff8cf1baca266b760d6a8dfe Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Thu, 24 Jul 2025 18:23:03 -0700 Subject: [PATCH] implement basic stream problem detection --- hack/upload-fixture.sh | 0 js/app/src/screens/live-dashboard.tsx | 20 +++- .../src/livestream-store/livestream-state.tsx | 9 ++ .../src/livestream-store/livestream-store.tsx | 2 + .../src/livestream-store/problems.tsx | 96 +++++++++++++++++++ .../livestream-store/websocket-consumer.tsx | 10 ++ .../docs/guides/start-streaming/obs.md | 2 +- .../lex-reference/place-stream-segment.md | 4 + lexicons/place/stream/segment.json | 3 +- pkg/media/audio_smear.go | 10 +- pkg/media/media_data_parser.go | 57 ++++++++++- pkg/media/media_data_parser_test.go | 19 ++++ pkg/model/segment.go | 10 +- pkg/streamplace/cbor_gen.go | 64 ++++++++++++- pkg/streamplace/streamsegment.go | 1 + 15 files changed, 293 insertions(+), 14 deletions(-) mode change 100644 => 100755 hack/upload-fixture.sh create mode 100644 js/components/src/livestream-store/problems.tsx diff --git a/hack/upload-fixture.sh b/hack/upload-fixture.sh old mode 100644 new mode 100755 diff --git a/js/app/src/screens/live-dashboard.tsx b/js/app/src/screens/live-dashboard.tsx index 996119eb..9e98c941 100644 --- a/js/app/src/screens/live-dashboard.tsx +++ b/js/app/src/screens/live-dashboard.tsx @@ -1,4 +1,7 @@ -import { LivestreamProvider } from "@streamplace/components"; +import { + LivestreamProvider, + useLivestreamStore, +} from "@streamplace/components"; import { Camera, FerrisWheel, X } from "@tamagui/lucide-icons"; import { Redirect } from "components/aqlink"; import CreateLivestream from "components/create-livestream"; @@ -134,6 +137,9 @@ export default function LiveDashboard() { {page === "update" && isLive ? : null} {page === "create" ? : null} + + + {madeChoiceAboutDebugRecording ? null : } @@ -141,6 +147,18 @@ export default function LiveDashboard() { ); } +const Problems = () => { + const problems = useLivestreamStore((x) => x.problems); + if (problems.length === 0) { + return null; + } + return ( + + {JSON.stringify(problems, null, 2)} + + ); +}; + const elems = [ { title: "Stream your camera!", diff --git a/js/components/src/livestream-store/livestream-state.tsx b/js/components/src/livestream-store/livestream-state.tsx index cf6a8b76..bdefb120 100644 --- a/js/components/src/livestream-store/livestream-state.tsx +++ b/js/components/src/livestream-store/livestream-state.tsx @@ -15,8 +15,17 @@ export interface LivestreamState { viewers: number | null; pendingHides: string[]; segment: PlaceStreamSegment.Record | null; + recentSegments: PlaceStreamSegment.Record[]; + problems: LivestreamProblem[]; renditions: PlaceStreamDefs.Rendition[]; replyToMessage: ChatMessageViewHydrated | null; streamKey: string | null; setStreamKey: (key: string | null) => void; } + +export interface LivestreamProblem { + code: string; + message: string; + severity: "error" | "warning" | "info"; + link?: string; +} diff --git a/js/components/src/livestream-store/livestream-store.tsx b/js/components/src/livestream-store/livestream-store.tsx index c73f8b8b..a5fb31e8 100644 --- a/js/components/src/livestream-store/livestream-store.tsx +++ b/js/components/src/livestream-store/livestream-store.tsx @@ -20,6 +20,8 @@ export const makeLivestreamStore = (): StoreApi => { streamKey: null, setStreamKey: (sk) => set({ streamKey: sk }), authors: {}, + recentSegments: [], + problems: [], })); }; diff --git a/js/components/src/livestream-store/problems.tsx b/js/components/src/livestream-store/problems.tsx new file mode 100644 index 00000000..bcce311e --- /dev/null +++ b/js/components/src/livestream-store/problems.tsx @@ -0,0 +1,96 @@ +import { PlaceStreamSegment } from "streamplace"; +import { LivestreamProblem } from "./livestream-state"; + +const VARIANCE_THRESHOLD = 0.5; +const DURATION_THRESHOLD = 5000000000; // 5s in ns + +const detectVariableSegmentLength = ( + segments: PlaceStreamSegment.Record[], +): { variable: boolean; duration: boolean } => { + if (segments.length < 3) { + // Need at least 3 segments to detect variability + return { variable: false, duration: false }; + } + + const durations = segments + .map((segment) => segment.duration) + .filter( + (duration): duration is number => duration !== undefined && duration > 0, + ); + + if (durations.length < 3) { + return { variable: false, duration: false }; + } + + // Calculate mean + const mean = + durations.reduce((sum: number, duration: number) => sum + duration, 0) / + durations.length; + + // Calculate standard deviation + const variance = + durations.reduce((sum: number, duration: number) => { + const diff = duration - mean; + return sum + diff * diff; + }, 0) / durations.length; + const stdDev = Math.sqrt(variance); + + // Calculate coefficient of variation (CV) + const cv = stdDev / mean; + + // CV > 0.5 indicates high variability + // This threshold can be adjusted based on testing + return { + variable: cv > VARIANCE_THRESHOLD, + duration: mean > DURATION_THRESHOLD, + }; +}; + +export const findProblems = ( + segments: PlaceStreamSegment.Record[], +): LivestreamProblem[] => { + const problems: LivestreamProblem[] = []; + let hasBFrames = false; + for (const segment of segments) { + const video = segment.video?.[0]; + if (!video) { + // i mean yes this is a problem but it can't happen yet + continue; + } + if (video.bframes === true) { + hasBFrames = true; + break; + } + } + if (hasBFrames) { + problems.push({ + code: "bframes", + message: + "Your stream contains B-Frames, which are not supported in Streamplace. Your stream will stutter.", + severity: "error", + link: "https://stream.place/docs/guides/start-streaming/obs/#obs-configuration", + }); + } + + const { variable, duration } = detectVariableSegmentLength(segments); + if (variable) { + problems.push({ + code: "variable_segment_length", + message: + "Your stream contains variable segment lengths, which may cause playback issues.", + severity: "warning", + link: "https://stream.place/docs/guides/start-streaming/obs/#obs-configuration", + }); + } + if (duration) { + problems.push({ + code: "long_segments", + message: + "Your stream contains long segments (>5s). This will work fine, but increases the delay of the livestream.", + severity: "warning", + link: "https://stream.place/docs/guides/start-streaming/obs/#obs-configuration", + }); + } + + return problems; +}; diff --git a/js/components/src/livestream-store/websocket-consumer.tsx b/js/components/src/livestream-store/websocket-consumer.tsx index 54f41301..535dbb08 100644 --- a/js/components/src/livestream-store/websocket-consumer.tsx +++ b/js/components/src/livestream-store/websocket-consumer.tsx @@ -11,6 +11,9 @@ import { } from "streamplace"; import { reduceChat } from "./chat"; import { LivestreamState } from "./livestream-state"; +import { findProblems } from "./problems"; + +const MAX_RECENT_SEGMENTS = 10; export const handleWebSocketMessages = ( state: LivestreamState, @@ -40,9 +43,16 @@ export const handleWebSocketMessages = ( }; state = reduceChat(state, [hydrated], [], []); } else if (PlaceStreamSegment.isRecord(message)) { + const newRecentSegments = [...state.recentSegments]; + newRecentSegments.unshift(message); + if (newRecentSegments.length > MAX_RECENT_SEGMENTS) { + newRecentSegments.pop(); + } state = { ...state, segment: message as PlaceStreamSegment.Record, + recentSegments: newRecentSegments, + problems: findProblems(newRecentSegments), }; } else if (PlaceStreamDefs.isBlockView(message)) { const block = message as PlaceStreamDefs.BlockView; diff --git a/js/docs/src/content/docs/guides/start-streaming/obs.md b/js/docs/src/content/docs/guides/start-streaming/obs.md index 2bda423a..ff4792f1 100644 --- a/js/docs/src/content/docs/guides/start-streaming/obs.md +++ b/js/docs/src/content/docs/guides/start-streaming/obs.md @@ -26,7 +26,7 @@ sidebar: 6. Click "Generate Stream Key" - The stream key will automatically be copied to your clipboard -### 2. Configure OBS Studio +### 2. Configure OBS Studio #### 2a. Initial OBS Configuration diff --git a/js/docs/src/content/docs/lex-reference/place-stream-segment.md b/js/docs/src/content/docs/lex-reference/place-stream-segment.md index 27eb66b3..9efa2d34 100644 --- a/js/docs/src/content/docs/lex-reference/place-stream-segment.md +++ b/js/docs/src/content/docs/lex-reference/place-stream-segment.md @@ -61,6 +61,7 @@ Media file representing a segment of a livestream | `width` | `integer` | ✅ | | | | `height` | `integer` | ✅ | | | | `framerate` | [`#framerate`](#framerate) | ❌ | | | +| `bframes` | `boolean` | ❌ | | | --- @@ -180,6 +181,9 @@ Media file representing a segment of a livestream "framerate": { "type": "ref", "ref": "#framerate" + }, + "bframes": { + "type": "boolean" } } }, diff --git a/lexicons/place/stream/segment.json b/lexicons/place/stream/segment.json index 6dac3903..5c1652f4 100644 --- a/lexicons/place/stream/segment.json +++ b/lexicons/place/stream/segment.json @@ -67,7 +67,8 @@ "framerate": { "type": "ref", "ref": "#framerate" - } + }, + "bframes": { "type": "boolean" } } }, "framerate": { diff --git a/pkg/media/audio_smear.go b/pkg/media/audio_smear.go index 78b7d8c9..78b22d17 100644 --- a/pkg/media/audio_smear.go +++ b/pkg/media/audio_smear.go @@ -137,6 +137,11 @@ func ToBuffers(ctx context.Context, input io.Reader) (*SegmentData, error) { NeedDataFunc: ReaderNeedData(ctx, input), }) + seg := SegmentData{ + Audio: []SegmentBuffer{}, + Video: []SegmentBuffer{}, + } + audioSinkElem, err := pipeline.GetElementByName("audioappsink") if err != nil { return nil, fmt.Errorf("failed to get audioappsink element: %w", err) @@ -146,11 +151,6 @@ func ToBuffers(ctx context.Context, input io.Reader) (*SegmentData, error) { return nil, fmt.Errorf("failed to get audioappsink element: %w", err) } - seg := SegmentData{ - Audio: []SegmentBuffer{}, - Video: []SegmentBuffer{}, - } - audioSink.SetCallbacks(&app.SinkCallbacks{ NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { sample := sink.PullSample() diff --git a/pkg/media/media_data_parser.go b/pkg/media/media_data_parser.go index c69c1064..bbda3a22 100644 --- a/pkg/media/media_data_parser.go +++ b/pkg/media/media_data_parser.go @@ -21,7 +21,9 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*model.SegmentMed ctx, cancel := context.WithCancel(ctx) defer cancel() pipelineSlice := []string{ - "appsrc name=appsrc ! qtdemux name=demux ! fakesink sync=false", + "appsrc name=appsrc ! qtdemux name=demux", + "demux.video_0 ! queue ! h264parse name=videoparse disable-passthrough=true config-interval=-1 ! h264timestamper ! appsink sync=false name=videoappsink", + "demux.audio_0 ! queue ! opusparse name=audioparse ! appsink sync=false name=audioappsink", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) @@ -118,6 +120,57 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*model.SegmentMed return nil, fmt.Errorf("error connecting pad-add: %w", err) } + audioSinkElem, err := pipeline.GetElementByName("audioappsink") + if err != nil { + return nil, fmt.Errorf("failed to get audioappsink element: %w", err) + } + audioSink := app.SinkFromElement(audioSinkElem) + if audioSink == nil { + return nil, fmt.Errorf("failed to get audioappsink element: %w", err) + } + + audioSink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowOK + } + + return gst.FlowOK + }, + }) + + videoSinkElem, err := pipeline.GetElementByName("videoappsink") + if err != nil { + return nil, fmt.Errorf("failed to get videoappsink element: %w", err) + } + videoSink := app.SinkFromElement(videoSinkElem) + if videoSink == nil { + return nil, fmt.Errorf("failed to get videoappsink element: %w", err) + } + + hasBFrames := false + videoSink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowOK + } + + buf := sample.GetBuffer() + pts := buf.PresentationTimestamp().String() + dts := buf.DecodingTimestamp().String() + + if pts != dts { + hasBFrames = true + } else { + log.Log(ctx, "no bframes", "pts", pts, "dts", dts) + } + + return gst.FlowOK + }, + }) + go func() { if err := HandleBusMessages(ctx, pipeline); err != nil { log.Log(ctx, "pipeline error", "error", err) @@ -145,6 +198,8 @@ func ParseSegmentMediaData(ctx context.Context, mp4bs []byte) (*model.SegmentMed return nil, fmt.Errorf("no audio metadata") } + videoMetadata.BFrames = hasBFrames + meta := &model.SegmentMediaData{ Video: []*model.SegmentMediadataVideo{videoMetadata}, Audio: []*model.SegmentMediadataAudio{audioMetadata}, diff --git a/pkg/media/media_data_parser_test.go b/pkg/media/media_data_parser_test.go index c5333f8d..bfbe6e8b 100644 --- a/pkg/media/media_data_parser_test.go +++ b/pkg/media/media_data_parser_test.go @@ -8,6 +8,7 @@ import ( "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/test/remote" ) func TestMediaDataParser(t *testing.T) { @@ -23,6 +24,24 @@ func TestMediaDataParser(t *testing.T) { mediaData, err := ParseSegmentMediaData(ctx, bs) require.NoError(t, err) require.NotNil(t, mediaData) + require.False(t, mediaData.Video[0].BFrames, "Video should not have BFrames") + require.Greater(t, mediaData.Duration, int64(0), "Video duration should not be empty") + }) +} + +func TestMediaDataParserBFrames(t *testing.T) { + withNoGSTLeaks(t, func() { + inputFile, err := os.Open(remote.RemoteFixture("5ea6c4491bade0cdcad3770aa0b63b2cd7a580e233ee320d5bc2282503b26491/segment-with-bframes.mp4")) + require.NoError(t, err) + defer inputFile.Close() + bs, err := io.ReadAll(inputFile) + require.NoError(t, err) + + ctx := log.WithDebugValue(context.Background(), map[string]map[string]int{"GStreamerFunc": {"ParseSegmentMediaData": 9}}) + mediaData, err := ParseSegmentMediaData(ctx, bs) + require.NoError(t, err) + require.NotNil(t, mediaData) + require.True(t, mediaData.Video[0].BFrames, "Video should have BFrames") require.Greater(t, mediaData.Duration, int64(0), "Video duration should not be empty") }) } diff --git a/pkg/model/segment.go b/pkg/model/segment.go index 035820b1..8f7afa03 100644 --- a/pkg/model/segment.go +++ b/pkg/model/segment.go @@ -15,10 +15,11 @@ import ( ) type SegmentMediadataVideo struct { - Width int `json:"width"` - Height int `json:"height"` - FPSNum int `json:"fpsNum"` - FPSDen int `json:"fpsDen"` + Width int `json:"width"` + Height int `json:"height"` + FPSNum int `json:"fpsNum"` + FPSDen int `json:"fpsDen"` + BFrames bool `json:"bframes"` } type SegmentMediadataAudio struct { @@ -89,6 +90,7 @@ func (s *Segment) ToStreamplaceSegment() (*streamplace.Segment, error) { Num: int64(s.MediaData.Video[0].FPSNum), Den: int64(s.MediaData.Video[0].FPSDen), }, + Bframes: &s.MediaData.Video[0].BFrames, }, }, Audio: []*streamplace.Segment_Audio{ diff --git a/pkg/streamplace/cbor_gen.go b/pkg/streamplace/cbor_gen.go index e0b19e48..a5607ef3 100644 --- a/pkg/streamplace/cbor_gen.go +++ b/pkg/streamplace/cbor_gen.go @@ -1224,7 +1224,11 @@ func (t *Segment_Video) MarshalCBOR(w io.Writer) error { } cw := cbg.NewCborWriter(w) - fieldCount := 4 + fieldCount := 5 + + if t.Bframes == nil { + fieldCount-- + } if t.Framerate == nil { fieldCount-- @@ -1301,6 +1305,31 @@ func (t *Segment_Video) MarshalCBOR(w io.Writer) error { } } + // t.Bframes (bool) (bool) + if t.Bframes != nil { + + if len("bframes") > 1000000 { + return xerrors.Errorf("Value in field \"bframes\" was too long") + } + + if err := cw.WriteMajorTypeHeader(cbg.MajTextString, uint64(len("bframes"))); err != nil { + return err + } + if _, err := cw.WriteString(string("bframes")); err != nil { + return err + } + + if t.Bframes == nil { + if _, err := cw.Write(cbg.CborNull); err != nil { + return err + } + } else { + if err := cbg.WriteBool(w, *t.Bframes); err != nil { + return err + } + } + } + // t.Framerate (streamplace.Segment_Framerate) (struct) if t.Framerate != nil { @@ -1426,6 +1455,39 @@ func (t *Segment_Video) UnmarshalCBOR(r io.Reader) (err error) { t.Height = int64(extraI) } + // t.Bframes (bool) (bool) + case "bframes": + + { + b, err := cr.ReadByte() + if err != nil { + return err + } + if b != cbg.CborNull[0] { + if err := cr.UnreadByte(); err != nil { + return err + } + + maj, extra, err = cr.ReadHeader() + if err != nil { + return err + } + if maj != cbg.MajOther { + return fmt.Errorf("booleans must be major type 7") + } + + var val bool + switch extra { + case 20: + val = false + case 21: + val = true + default: + return fmt.Errorf("booleans are either major type 7, value 20 or 21 (got %d)", extra) + } + t.Bframes = &val + } + } // t.Framerate (streamplace.Segment_Framerate) (struct) case "framerate": diff --git a/pkg/streamplace/streamsegment.go b/pkg/streamplace/streamsegment.go index eecb6b36..a97cf82d 100644 --- a/pkg/streamplace/streamsegment.go +++ b/pkg/streamplace/streamsegment.go @@ -48,6 +48,7 @@ type Segment_SegmentView struct { // Segment_Video is a "video" in the place.stream.segment schema. type Segment_Video struct { + Bframes *bool `json:"bframes,omitempty" cborgen:"bframes,omitempty"` Codec string `json:"codec" cborgen:"codec"` Framerate *Segment_Framerate `json:"framerate,omitempty" cborgen:"framerate,omitempty"` Height int64 `json:"height" cborgen:"height"` -- 2.51.2