test(gateway): measure qualification wire capacity
Verify Data Plane / gateway (push) Successful in 4m46s

This commit is contained in:
sechmachine
2026-08-09 23:01:42 +07:00
parent 122080ab34
commit 22433e5c45
7 changed files with 176 additions and 24 deletions
+9 -1
View File
@@ -940,12 +940,20 @@ func (c *independentGatewayClient) ReceiveFrame(ctx context.Context) (Frame, err
}
func (c *independentGatewayClient) ReceiveMedia(ctx context.Context) ([]byte, error) {
return c.receiveMedia(ctx, nil)
}
func (c *independentGatewayClient) receiveMedia(ctx context.Context, observe func(time.Time, int)) ([]byte, error) {
for {
data, err := c.connection.ReceiveDatagram(ctx)
if err != nil {
return nil, err
}
payload, complete, err := c.media.Add(data, time.Now())
receivedAt := time.Now()
if observe != nil {
observe(receivedAt, len(data))
}
payload, complete, err := c.media.Add(data, receivedAt)
if err != nil {
return nil, err
}
+111
View File
@@ -669,6 +669,117 @@ func TestQualificationLossAndSteppedThroughputBounds(t *testing.T) {
profile.Name == "constrained" && len(observation.CapacityStepObservations) != 2 {
t.Fatalf("%s observation = %#v", profile.Name, observation)
}
for _, step := range observation.CapacityStepObservations {
t.Logf("%s %d%%: convergence=%s wire_max_5s=%d wire_cap=%d", profile.Name, step.ReductionPercent, step.Convergence, step.MaximumFiveSecond, step.FiveSecondCap)
}
}
}
func TestQualificationCapacityStepsMeasurePublicWireDatagrams(t *testing.T) {
const packetCount = qualificationImpairmentPacketCount
profile := qualificationMediaProfiles()[0]
started := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC)
spacing := time.Duration(int64(time.Second) * int64(profile.PacketBytes) * 8 / (profile.BitrateKbps * 1000))
stepAt := map[int]time.Time{
25: started.Add(time.Duration(packetCount/3) * spacing),
50: started.Add(time.Duration(packetCount*2/3) * spacing),
}
if phase := stepAt[50].Sub(stepAt[25]); phase != 1_571_842_800*time.Nanosecond {
t.Fatalf("25%% phase = %s, want exact source-shaped transition", phase)
}
pacer := newFairPacer(qualificationMediaPacerKbps(profile, 0))
const flow = "qualification-wire-flow"
now := started
stalled := false
quarterApplied := false
halfApplied := false
debtBeforeQuarter := time.Duration(0)
logicalDeliveries := make([]qualificationDeliverySample, 0, packetCount)
type wireSample struct {
reserved time.Time
delivery qualificationDeliverySample
}
wireSamples := make([]wireSample, 0, packetCount*2)
for index := 0; index < packetCount; index++ {
sourceAt := started.Add(time.Duration(index) * spacing)
if now.Before(sourceAt) {
now = sourceAt
}
if !stalled && !sourceAt.Before(stepAt[25].Add(-150*time.Millisecond)) {
now = now.Add(42 * time.Millisecond)
stalled = true
}
for _, encodedBytes := range []int{1200, 25} {
if !quarterApplied && !now.Before(stepAt[25]) {
pacer.mu.Lock()
debtBeforeQuarter = pacer.flows[flow].debt
pacer.mu.Unlock()
pacer.setKbps(qualificationMediaPacerKbps(profile, 25))
quarterApplied = true
}
if !halfApplied && !now.Before(stepAt[50]) {
pacer.setKbps(qualificationMediaPacerKbps(profile, 50))
halfApplied = true
}
reserved := pacer.reserveAt(now, flow, encodedBytes)
if reserved.After(now) {
now = reserved
}
wireSamples = append(wireSamples, wireSample{
reserved: reserved,
delivery: qualificationDeliverySample{At: now, Bytes: int64(encodedBytes)},
})
}
logicalDeliveries = append(logicalDeliveries, qualificationDeliverySample{At: now, Bytes: 1225})
}
if !quarterApplied || !halfApplied || debtBeforeQuarter <= 0 || debtBeforeQuarter > nativeApolloVideoQueueLatency-fairPacerMaximumCatchup {
t.Fatalf("capacity transition state quarter=%t half=%t debt=%s", quarterApplied, halfApplied, debtBeforeQuarter)
}
pacer.mu.Lock()
remainingDebt := pacer.flows[flow].debt
pacer.mu.Unlock()
if remainingDebt != 0 {
t.Fatalf("remaining debt = %s, want zero", remainingDebt)
}
t.Logf("wire model: datagrams=%d debt_before_25=%s remaining_debt=%s", len(wireSamples), debtBeforeQuarter, remainingDebt)
wireDeliveries := make([]qualificationDeliverySample, len(wireSamples))
for index, sample := range wireSamples {
if sample.delivery.At.Before(sample.reserved) {
t.Fatalf("wire datagram %d delivered at %s before reservation %s", index, sample.delivery.At, sample.reserved)
}
wireDeliveries[index] = sample.delivery
}
for _, reduction := range []int{25, 50} {
wireBytesPerSecond := qualificationMediaPacerKbps(profile, reduction) * 1000 / 8
wireAfterStep := qualificationDeliveriesAfter(wireDeliveries, stepAt[reduction])
maximum := qualificationMaximumDeliveryBytes(wireAfterStep, 5*time.Second)
if maximum > wireBytesPerSecond*5*105/100 {
t.Fatalf("%d%% wire five-second maximum = %d, cap = %d", reduction, maximum, wireBytesPerSecond*5*105/100)
}
payloadBytesPerSecond := profile.BitrateKbps * int64(100-reduction) * 1000 / 100 / 8
legacy := qualificationMeasuredConvergence(
qualificationDeliveriesAfter(logicalDeliveries, stepAt[reduction]), stepAt[reduction], payloadBytesPerSecond,
)
if reduction == 25 && legacy != 11*time.Second {
t.Fatalf("25%% complete-frame classifier convergence = %s, want 11s sentinel reproduction", legacy)
}
observation := qualificationCapacityStepObservation(wireDeliveries, stepAt[reduction], profile, reduction)
if observation.Convergence > 10*time.Second {
for offset := time.Duration(0); offset < 2*time.Second; offset += 250 * time.Millisecond {
var bucket int64
for _, delivery := range wireAfterStep {
if !delivery.At.Before(stepAt[reduction].Add(offset)) && delivery.At.Before(stepAt[reduction].Add(offset+250*time.Millisecond)) {
bucket += delivery.Bytes
}
}
t.Logf("%d%% wire bucket %s = %d bytes (%d B/s)", reduction, offset, bucket, bucket*4)
}
t.Fatalf("%d%% wire convergence = %s sentinel with wire maximum %d under cap %d", reduction, observation.Convergence, maximum, wireBytesPerSecond*5*105/100)
}
t.Logf("%d%%: payload_target=%d wire_target=%d legacy=%s wire_convergence=%s wire_max_5s=%d wire_cap_105=%d",
reduction, payloadBytesPerSecond, wireBytesPerSecond, legacy, observation.Convergence, maximum, wireBytesPerSecond*5*105/100)
}
}
+34 -20
View File
@@ -38,7 +38,7 @@ import (
)
const (
qualificationToolVersion = "versevdi-gateway-qualification/v8"
qualificationToolVersion = "versevdi-gateway-qualification/v9"
qualificationImpairmentQueuePackets = nativeApolloVideoQueuePackets
qualificationImpairmentMaxPackets = 100_000
qualificationImpairmentPacketCount = 10_000
@@ -1243,9 +1243,13 @@ func (p *qualificationPath) processingDiagnostics(processed int64) qualification
}
func (p *qualificationPath) receivePayload(parent context.Context) ([]byte, error) {
return p.receivePayloadObserved(parent, nil)
}
func (p *qualificationPath) receivePayloadObserved(parent context.Context, observe func(time.Time, int)) ([]byte, error) {
ctx, cancel := context.WithTimeout(parent, 2*time.Second)
defer cancel()
return p.client.ReceiveMedia(ctx)
return p.client.receiveMedia(ctx, observe)
}
func qualificationSourceVideoPackets(t *testing.T, key []byte, frame uint32, encoded []byte) [][]byte {
@@ -1503,6 +1507,7 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
beforeProviderDrops := path.session.mediaDrops.Load()
beforeSourceUDP := path.sourceUDP.Load()
beforePacer := path.server.pacer.reservations.Load()
wireDeliveries := make([]qualificationDeliverySample, 0, len(jobs)*2)
grace := max(2*profile.RTT+2*profile.Jitter, 2*time.Second)
lastTarget := time.Duration(packetCount) * spacing
if len(jobs) > 0 {
@@ -1522,7 +1527,9 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
metrics := beforeMetrics
seen := make([]bool, packetCount)
for len(result.packets) < len(jobs) {
recovered, receiveErr := path.receivePayload(receiveCtx)
recovered, receiveErr := path.receivePayloadObserved(receiveCtx, func(receivedAt time.Time, bytes int) {
wireDeliveries = append(wireDeliveries, qualificationDeliverySample{At: receivedAt, Bytes: int64(bytes)})
})
if receiveErr != nil {
if errors.Is(receiveErr, context.DeadlineExceeded) || errors.Is(receiveErr, context.Canceled) {
break
@@ -1595,7 +1602,6 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
return qualificationImpairmentObservation{}, received.err
}
var deliveries []qualificationDeliverySample
var totalLatency, totalJitter, previousLatency time.Duration
previousDelivered := -1
for _, packet := range received.packets {
@@ -1620,10 +1626,6 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
observation.ObservedOutOfOrder++
}
previousDelivered = packet.index
datagrams := (media.PacketBytes + frameV2PayloadSize - 1) / frameV2PayloadSize
deliveries = append(deliveries, qualificationDeliverySample{
At: packet.deliveredAt, Bytes: int64(media.PacketBytes + datagrams*frameV2HeaderSize),
})
observation.MaxQueuePackets = max(observation.MaxQueuePackets, packet.queuePackets)
}
observation.Delivered = len(received.packets)
@@ -1712,17 +1714,11 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
return qualificationImpairmentObservation{}, err
}
for _, reduction := range profile.CapacitySteps {
bytesPerSecond := media.BitrateKbps * int64(100-reduction) * 1000 / 100 / 8
stepDeliveries := qualificationDeliveriesAfter(deliveries, stepAt[reduction])
convergence := qualificationMeasuredConvergence(stepDeliveries, stepAt[reduction], bytesPerSecond)
maximum := qualificationMaximumDeliveryBytes(stepDeliveries, 5*time.Second)
observation.CapacityStepObservations = append(observation.CapacityStepObservations, qualificationCapacityStep{
ReductionPercent: reduction, Convergence: convergence,
MaximumFiveSecond: maximum, FiveSecondCap: bytesPerSecond * 5,
})
step := qualificationCapacityStepObservation(wireDeliveries, stepAt[reduction], media, reduction)
observation.CapacityStepObservations = append(observation.CapacityStepObservations, step)
if packetCount >= qualificationImpairmentPacketCount &&
(convergence > 10*time.Second || maximum > bytesPerSecond*5*105/100) {
return qualificationImpairmentObservation{}, fmt.Errorf("capacity step %d failed measured convergence=%s five-second=%d", reduction, convergence, maximum)
(step.Convergence > 10*time.Second || step.MaximumFiveSecond > step.FiveSecondCap*105/100) {
return qualificationImpairmentObservation{}, fmt.Errorf("capacity step %d failed measured convergence=%s five-second=%d", reduction, step.Convergence, step.MaximumFiveSecond)
}
}
if observation.Delivered+observation.Dropped != observation.Sent ||
@@ -1733,6 +1729,17 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
return observation, nil
}
func qualificationCapacityStepObservation(wireDeliveries []qualificationDeliverySample, start time.Time, media qualificationMediaProfile, reduction int) qualificationCapacityStep {
bytesPerSecond := qualificationMediaPacerKbps(media, reduction) * 1000 / 8
deliveries := qualificationDeliveriesAfter(wireDeliveries, start)
return qualificationCapacityStep{
ReductionPercent: reduction,
Convergence: qualificationMeasuredConvergence(deliveries, start, bytesPerSecond),
MaximumFiveSecond: qualificationMaximumDeliveryBytes(deliveries, 5*time.Second),
FiveSecondCap: bytesPerSecond * 5,
}
}
func qualificationKnownImpairment(profile qualificationImpairmentProfile) bool {
for _, known := range qualificationImpairmentProfiles() {
if profile.Name != known.Name || profile.RTT != known.RTT || profile.Jitter != known.Jitter ||
@@ -1759,9 +1766,16 @@ func qualificationDeliveriesAfter(deliveries []qualificationDeliverySample, star
func qualificationMeasuredConvergence(deliveries []qualificationDeliverySample, start time.Time, targetBytesPerSecond int64) time.Duration {
const window = 250 * time.Millisecond
const requiredWindows = 4
if len(deliveries) == 0 {
return 11 * time.Second
}
windowOrigin := start
if deliveries[0].At.After(windowOrigin) {
windowOrigin = deliveries[0].At
}
consecutive := 0
for offset := time.Duration(0); offset <= 10*time.Second; offset += window {
windowStart := start.Add(offset)
windowStart := windowOrigin.Add(offset)
var total int64
for _, delivery := range deliveries {
if !delivery.At.Before(windowStart) && delivery.At.Before(windowStart.Add(window)) {
@@ -1772,7 +1786,7 @@ func qualificationMeasuredConvergence(deliveries []qualificationDeliverySample,
if rate >= targetBytesPerSecond*90/100 && rate <= targetBytesPerSecond*105/100 {
consecutive++
if consecutive == requiredWindows {
return offset + window
return windowStart.Add(window).Sub(start)
}
} else {
consecutive = 0