fix(stream): isolate H264 subscribers

Keep slow WebRTC consumers from blocking the shared H264 source by dropping buffered frames until the next keyframe. Restore the WebRTC NACK and RTCP report interceptors while retaining a safe RTP MTU, and enable linker --as-needed for libkvm to remove unused OpenCV dependencies.
This commit is contained in:
watermeko
2026-08-28 09:49:45 +00:00
committed by Guoguo
parent a9e1ef2459
commit 7f95fe9bbd
6 changed files with 43 additions and 7 deletions

Binary file not shown.

Binary file not shown.

View File

@@ -15,10 +15,11 @@ type H264Frame struct {
} }
type H264Subscription struct { type H264Subscription struct {
frames chan H264Frame frames chan H264Frame
done chan struct{} done chan struct{}
closed atomic.Bool closed atomic.Bool
once sync.Once once sync.Once
waitingForKeyframe bool
} }
type H264Source struct { type H264Source struct {
@@ -147,11 +148,38 @@ func (s *H264Subscription) send(frame H264Frame) bool {
if s.closed.Load() { if s.closed.Load() {
return false return false
} }
if s.waitingForKeyframe {
if frame.Result != 3 {
return false
}
s.waitingForKeyframe = false
}
select { select {
case s.frames <- frame: case s.frames <- frame:
return true return true
case <-s.done: case <-s.done:
return false return false
default:
}
for {
select {
case <-s.frames:
case <-s.done:
return false
default:
s.waitingForKeyframe = true
if frame.Result != 3 {
return false
}
s.waitingForKeyframe = false
select {
case s.frames <- frame:
return true
case <-s.done:
return false
}
}
} }
} }

View File

@@ -164,9 +164,16 @@ func createPeerConnection(iceServers []webrtc.ICEServer, mediaEngine *webrtc.Med
apiOptions := []func(api *webrtc.API){ apiOptions := []func(api *webrtc.API){
webrtc.WithSettingEngine(settingEngine), webrtc.WithSettingEngine(settingEngine),
webrtc.WithInterceptorRegistry(&interceptor.Registry{}),
} }
if mediaEngine != nil { if mediaEngine != nil {
registry := &interceptor.Registry{}
if err := webrtc.ConfigureNack(mediaEngine, registry); err != nil {
return nil, err
}
if err := webrtc.ConfigureRTCPReports(registry); err != nil {
return nil, err
}
apiOptions = append(apiOptions, webrtc.WithInterceptorRegistry(registry))
apiOptions = append(apiOptions, webrtc.WithMediaEngine(mediaEngine)) apiOptions = append(apiOptions, webrtc.WithMediaEngine(mediaEngine))
} }

View File

@@ -15,7 +15,7 @@ func NewWebRTCManager() *WebRTCManager {
m := &WebRTCManager{ m := &WebRTCManager{
clients: make(map[*websocket.Conn]*Client), clients: make(map[*websocket.Conn]*Client),
videoPacketizer: rtp.NewPacketizer( videoPacketizer: rtp.NewPacketizer(
1450, 1200,
100, 100,
0x1234ABCD, 0x1234ABCD,
&codecs.H264Payloader{}, &codecs.H264Payloader{},

View File

@@ -1,5 +1,6 @@
list(APPEND ADD_REQUIREMENTS vision basic peripheral) list(APPEND ADD_REQUIREMENTS vision basic peripheral)
list(APPEND ADD_INCLUDE "include") list(APPEND ADD_INCLUDE "include")
list(APPEND ADD_LINK_DEFINITIONS_PRIVATE -Wl,--as-needed)
append_srcs_dir(ADD_SRCS "src") append_srcs_dir(ADD_SRCS "src")