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 996119ebf..9e98c9412 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 cf6a8b762..bdefb120a 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 c73f8b8bd..a5fb31e8e 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 000000000..bcce311ec
--- /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 54f41301a..535dbb086 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 2bda423a7..ff4792f16 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 27eb66b33..9efa2d346 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 6dac39033..5c1652f45 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 78b7d8c96..78b22d172 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 c69c10641..bbda3a224 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 c5333f8d1..bfbe6e8b6 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 035820b18..8f7afa03e 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 e0b19e486..a5607ef39 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 eecb6b367..a97cf82d4 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"`