fix(gateway): bound Apollo video ingress
Verify Data Plane / gateway (push) Successful in 4m14s

This commit is contained in:
sechmachine
2026-08-09 19:07:22 +07:00
parent a0ca194691
commit 55afea72a1
8 changed files with 467 additions and 39 deletions
+126 -36
View File
@@ -33,6 +33,9 @@ const (
nativeApolloVideoQueueLatency = 250 * time.Millisecond
nativeApolloAudioQueuePackets = 16
nativeApolloEventQueuePackets = 16
nativeApolloVideoIngressSlots = 2048
nativeApolloVideoPacketBytes = apolloVideoHeaderSize + apolloVideoRawPacketSize
nativeApolloVideoReadBuffer = nativeApolloVideoIngressSlots * nativeApolloVideoPacketBytes
)
// NativeApolloBackend keeps provider sockets inside the gateway process. The
@@ -41,8 +44,9 @@ const (
type NativeApolloBackend struct {
Dialer *net.Dialer
mu sync.Mutex
pending map[string]*apolloRTSPSetup
mu sync.Mutex
pending map[string]*apolloRTSPSetup
configureMedia func(*apolloMediaCodec)
}
func NewNativeApolloBackend() *NativeApolloBackend {
@@ -268,7 +272,7 @@ func (b *NativeApolloBackend) Open(ctx context.Context, request LaunchRequest, _
if !ok || setup == nil {
return nil, ErrProviderDisconnected
}
return newNativeApolloProviderSession(ctx, setup)
return newNativeApolloProviderSession(ctx, setup, b.configureMedia)
}
func readBounded(reader io.Reader, max int) ([]byte, error) {
@@ -336,7 +340,7 @@ func newNativeApolloSession(sessionID string) *nativeApolloSession {
}
}
func newNativeApolloProviderSession(ctx context.Context, setup *apolloRTSPSetup) (*nativeApolloSession, error) {
func newNativeApolloProviderSession(ctx context.Context, setup *apolloRTSPSetup, configureMedia func(*apolloMediaCodec)) (*nativeApolloSession, error) {
if setup == nil || len(setup.streamKey) != 16 || setup.controlPort == 0 || setup.audioPort == 0 || setup.videoPort == 0 {
return nil, ErrProviderMalformed
}
@@ -360,6 +364,9 @@ func newNativeApolloProviderSession(ctx context.Context, setup *apolloRTSPSetup)
peer.close(err)
return nil, err
}
if configureMedia != nil {
configureMedia(media)
}
session := newNativeApolloSession(setup.sessionID)
managementClient, err := newPinnedApolloHTTPClient(setup.providerWork)
if err != nil {
@@ -391,6 +398,12 @@ func newNativeApolloProviderSession(ctx context.Context, setup *apolloRTSPSetup)
peer.close(err)
return nil, err
}
if err := videoConn.SetReadBuffer(nativeApolloVideoReadBuffer); err != nil {
_ = audioConn.Close()
_ = videoConn.Close()
peer.close(err)
return nil, fmt.Errorf("set Apollo video receive buffer: %w", err)
}
if _, err := audioConn.Write(apolloMediaPing(setup.audioPing, 1)); err != nil {
_ = audioConn.Close()
_ = videoConn.Close()
@@ -896,19 +909,20 @@ func (s *nativeApolloSession) readUDPMedia() {
s.closeMediaChannels()
return
}
var readers sync.WaitGroup
readers.Add(2)
read := func(conn *net.UDPConn, output chan ProviderMedia, video bool) {
defer readers.Done()
videoIngress := newNativeApolloVideoIngress()
var workers sync.WaitGroup
workers.Add(3)
go func() {
defer workers.Done()
buffer := make([]byte, apolloMediaMaximumPacket+1)
for {
if s.mediaQuiesced.Load() {
return
}
if err := conn.SetReadDeadline(time.Now().Add(250 * time.Millisecond)); err != nil {
if err := s.audioConn.SetReadDeadline(time.Now().Add(250 * time.Millisecond)); err != nil {
return
}
count, err := conn.Read(buffer)
count, err := s.audioConn.Read(buffer)
receivedAt := time.Now()
if err != nil {
if networkErr, ok := err.(net.Error); ok && networkErr.Timeout() {
@@ -928,45 +942,121 @@ func (s *nativeApolloSession) readUDPMedia() {
return
}
s.mediaIngress.Add(1)
var payloads [][]byte
if video {
shard, openErr := s.media.OpenVideo(buffer[:count])
if openErr != nil {
continue
}
payload, err := s.videoFEC.Add(shard)
if err != nil || len(payload) == 0 {
continue
}
payloads = [][]byte{payload}
} else {
shard, openErr := s.media.OpenAudio(buffer[:count])
if openErr != nil {
continue
}
var evicted bool
payloads, evicted, err = s.audioFEC.Add(s.media, shard)
if evicted {
s.mediaDrops.Add(1)
}
shard, openErr := s.media.OpenAudio(buffer[:count])
if openErr != nil {
continue
}
payloads, evicted, err := s.audioFEC.Add(s.media, shard)
if evicted {
s.mediaDrops.Add(1)
}
if err != nil {
continue
}
for _, payload := range payloads {
s.enqueueMedia(output, payload, receivedAt)
s.enqueueMedia(s.audio, payload, receivedAt)
}
}
}
go read(s.audioConn, s.audio, false)
go read(s.videoConn, s.video, true)
}()
go func() {
readers.Wait()
defer workers.Done()
s.drainApolloVideo(videoIngress)
}()
go func() {
defer workers.Done()
s.processApolloVideo(videoIngress)
}()
go func() {
workers.Wait()
close(s.readDone)
s.closeMediaChannels()
}()
}
type nativeApolloVideoIngressSlot struct {
packet [apolloMediaMaximumPacket + 1]byte
count int
receivedAt time.Time
}
type nativeApolloVideoIngress struct {
slots [nativeApolloVideoIngressSlots]nativeApolloVideoIngressSlot
free chan uint16
ready chan uint16
scratch [apolloMediaMaximumPacket + 1]byte
}
func newNativeApolloVideoIngress() *nativeApolloVideoIngress {
ingress := &nativeApolloVideoIngress{
free: make(chan uint16, nativeApolloVideoIngressSlots),
ready: make(chan uint16, nativeApolloVideoIngressSlots),
}
for index := range nativeApolloVideoIngressSlots {
ingress.free <- uint16(index)
}
return ingress
}
func (s *nativeApolloSession) drainApolloVideo(ingress *nativeApolloVideoIngress) {
defer close(ingress.ready)
for {
if s.mediaQuiesced.Load() {
return
}
select {
case index := <-ingress.free:
slot := &ingress.slots[index]
count, err := s.videoConn.Read(slot.packet[:])
receivedAt := time.Now()
if err != nil {
ingress.free <- index
return
}
if count > apolloMediaMaximumPacket {
ingress.free <- index
continue
}
if s.mediaQuiesced.Load() {
ingress.free <- index
return
}
s.mediaIngress.Add(1)
slot.count, slot.receivedAt = count, receivedAt
ingress.ready <- index
default:
count, err := s.videoConn.Read(ingress.scratch[:])
if err != nil {
return
}
if count > apolloMediaMaximumPacket {
continue
}
if s.mediaQuiesced.Load() {
return
}
s.mediaIngress.Add(1)
s.mediaDrops.Add(1)
}
}
}
func (s *nativeApolloSession) processApolloVideo(ingress *nativeApolloVideoIngress) {
for index := range ingress.ready {
slot := &ingress.slots[index]
if !s.mediaQuiesced.Load() {
shard, err := s.media.OpenVideo(slot.packet[:slot.count])
if err == nil {
payload, fecErr := s.videoFEC.Add(shard)
if fecErr == nil && len(payload) > 0 {
s.enqueueMedia(s.video, payload, slot.receivedAt)
}
}
}
slot.count, slot.receivedAt = 0, time.Time{}
ingress.free <- index
}
}
func pushLatest[T any](channel chan T, payload T) bool {
select {
case channel <- payload: