feat(gateway): repair native Apollo provider path
This commit is contained in:
+180
-15
@@ -2,33 +2,114 @@ package gateway
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Metrics struct {
|
||||
ActiveSessions atomic.Int64
|
||||
AdmissionRejects atomic.Uint64
|
||||
MediaDrops atomic.Uint64
|
||||
ProviderErrors atomic.Uint64
|
||||
InputRejected atomic.Uint64
|
||||
ActiveSessions atomic.Int64
|
||||
AdmittedSessions atomic.Uint64
|
||||
AdmissionRejects atomic.Uint64
|
||||
Reconnects atomic.Uint64
|
||||
DrainTransitions atomic.Uint64
|
||||
MediaDrops atomic.Uint64
|
||||
MediaPackets atomic.Uint64
|
||||
MediaBytes atomic.Uint64
|
||||
QueueDelayNanos atomic.Uint64
|
||||
ProcessingDelayNanos atomic.Uint64
|
||||
ProcessingSamples atomic.Uint64
|
||||
PacingDelayNanos atomic.Uint64
|
||||
ProviderErrors atomic.Uint64
|
||||
InputRejected atomic.Uint64
|
||||
ControlRTTNanos atomic.Uint64
|
||||
ControlJitterNanos atomic.Uint64
|
||||
ControlLossPPM atomic.Uint64
|
||||
PendingReliable atomic.Uint64
|
||||
ProviderState atomic.Uint64
|
||||
}
|
||||
|
||||
type MetricsSnapshot struct {
|
||||
ActiveSessions int64
|
||||
AdmissionRejects uint64
|
||||
MediaDrops uint64
|
||||
ProviderErrors uint64
|
||||
InputRejected uint64
|
||||
ActiveSessions int64
|
||||
AdmittedSessions uint64
|
||||
AdmissionRejects uint64
|
||||
Reconnects uint64
|
||||
DrainTransitions uint64
|
||||
MediaDrops uint64
|
||||
MediaPackets uint64
|
||||
MediaBytes uint64
|
||||
QueueDelayNanos uint64
|
||||
ProcessingDelayNanos uint64
|
||||
ProcessingSamples uint64
|
||||
PacingDelayNanos uint64
|
||||
ProviderErrors uint64
|
||||
InputRejected uint64
|
||||
ControlRTTNanos uint64
|
||||
ControlJitterNanos uint64
|
||||
ControlLossPPM uint64
|
||||
PendingReliable uint64
|
||||
ProviderState uint64
|
||||
}
|
||||
|
||||
func (m *Metrics) Snapshot() MetricsSnapshot {
|
||||
return MetricsSnapshot{
|
||||
ActiveSessions: m.ActiveSessions.Load(),
|
||||
AdmissionRejects: m.AdmissionRejects.Load(),
|
||||
MediaDrops: m.MediaDrops.Load(),
|
||||
ProviderErrors: m.ProviderErrors.Load(),
|
||||
InputRejected: m.InputRejected.Load(),
|
||||
ActiveSessions: m.ActiveSessions.Load(),
|
||||
AdmittedSessions: m.AdmittedSessions.Load(),
|
||||
AdmissionRejects: m.AdmissionRejects.Load(),
|
||||
Reconnects: m.Reconnects.Load(),
|
||||
DrainTransitions: m.DrainTransitions.Load(),
|
||||
MediaDrops: m.MediaDrops.Load(),
|
||||
MediaPackets: m.MediaPackets.Load(),
|
||||
MediaBytes: m.MediaBytes.Load(),
|
||||
QueueDelayNanos: m.QueueDelayNanos.Load(),
|
||||
ProcessingDelayNanos: m.ProcessingDelayNanos.Load(),
|
||||
ProcessingSamples: m.ProcessingSamples.Load(),
|
||||
PacingDelayNanos: m.PacingDelayNanos.Load(),
|
||||
ProviderErrors: m.ProviderErrors.Load(),
|
||||
InputRejected: m.InputRejected.Load(),
|
||||
ControlRTTNanos: m.ControlRTTNanos.Load(),
|
||||
ControlJitterNanos: m.ControlJitterNanos.Load(),
|
||||
ControlLossPPM: m.ControlLossPPM.Load(),
|
||||
PendingReliable: m.PendingReliable.Load(),
|
||||
ProviderState: m.ProviderState.Load(),
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Metrics) observeProviderTelemetry(telemetry ProviderTelemetry) {
|
||||
if m == nil {
|
||||
return
|
||||
}
|
||||
m.ControlRTTNanos.Store(uint64(telemetry.ControlRTT))
|
||||
m.ControlJitterNanos.Store(uint64(telemetry.ControlJitter))
|
||||
m.PendingReliable.Store(telemetry.PendingReliable)
|
||||
if telemetry.ReliableSent == 0 {
|
||||
m.ControlLossPPM.Store(0)
|
||||
} else {
|
||||
m.ControlLossPPM.Store(telemetry.ReliableRetransmits * 1_000_000 / telemetry.ReliableSent)
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Metrics) observeProviderState(state string) {
|
||||
if m == nil {
|
||||
return
|
||||
}
|
||||
switch state {
|
||||
case ProviderStateStarting:
|
||||
m.ProviderState.Store(1)
|
||||
case ProviderStateReady:
|
||||
m.ProviderState.Store(2)
|
||||
case ProviderStateDisconnected:
|
||||
m.ProviderState.Store(3)
|
||||
case ProviderStateTerminating:
|
||||
m.ProviderState.Store(4)
|
||||
case ProviderStateTerminated:
|
||||
m.ProviderState.Store(5)
|
||||
case ProviderStateCleanup:
|
||||
m.ProviderState.Store(6)
|
||||
case ProviderStateFailed:
|
||||
m.ProviderState.Store(7)
|
||||
default:
|
||||
m.ProviderState.Store(0)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,6 +118,90 @@ type Pacer struct {
|
||||
last time.Time
|
||||
}
|
||||
|
||||
// fairPacer is the gateway's one shared, equal-tier media scheduler. Each
|
||||
// session can hold only its existing bounded provider media channel while it
|
||||
// waits for the next reservation, so a slow client cannot grow a global queue.
|
||||
type fairPacer struct {
|
||||
mu sync.Mutex
|
||||
bytesPerSecond int64
|
||||
flows map[string]fairPacerFlow
|
||||
}
|
||||
|
||||
type fairPacerFlow struct {
|
||||
next time.Time
|
||||
lastSeen time.Time
|
||||
}
|
||||
|
||||
func newFairPacer(kbps int64) *fairPacer {
|
||||
pacer := &fairPacer{flows: make(map[string]fairPacerFlow)}
|
||||
pacer.setKbps(kbps)
|
||||
return pacer
|
||||
}
|
||||
|
||||
func (p *fairPacer) setKbps(kbps int64) {
|
||||
if p == nil {
|
||||
return
|
||||
}
|
||||
p.mu.Lock()
|
||||
if kbps > 0 {
|
||||
p.bytesPerSecond = kbps * 1000 / 8
|
||||
} else {
|
||||
p.bytesPerSecond = 0
|
||||
}
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
func (p *fairPacer) remove(flow string) {
|
||||
if p == nil || flow == "" {
|
||||
return
|
||||
}
|
||||
p.mu.Lock()
|
||||
delete(p.flows, flow)
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
func (p *fairPacer) reserveAt(now time.Time, flow string, bytes int) time.Time {
|
||||
if p == nil || flow == "" || bytes < 1 {
|
||||
return now
|
||||
}
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
if p.bytesPerSecond < 1 {
|
||||
return now
|
||||
}
|
||||
for key, state := range p.flows {
|
||||
if now.Sub(state.lastSeen) > time.Second {
|
||||
delete(p.flows, key)
|
||||
}
|
||||
}
|
||||
state := p.flows[flow]
|
||||
state.lastSeen = now
|
||||
p.flows[flow] = state
|
||||
base := now
|
||||
if state.next.After(base) {
|
||||
base = state.next
|
||||
}
|
||||
numerator := int64(bytes) * int64(len(p.flows)) * int64(time.Second)
|
||||
delay := time.Duration((numerator + p.bytesPerSecond - 1) / p.bytesPerSecond)
|
||||
state.next = base.Add(delay)
|
||||
p.flows[flow] = state
|
||||
return state.next
|
||||
}
|
||||
|
||||
func (p *fairPacer) wait(ctx context.Context, flow string, bytes int) error {
|
||||
target := p.reserveAt(time.Now(), flow, bytes)
|
||||
if delay := time.Until(target); delay > 0 {
|
||||
timer := time.NewTimer(delay)
|
||||
defer timer.Stop()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-timer.C:
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func NewPacer(kbps int64) *Pacer {
|
||||
if kbps < 1 {
|
||||
return &Pacer{}
|
||||
|
||||
Reference in New Issue
Block a user