package gateway import ( "bufio" "bytes" "compress/gzip" "context" "crypto/aes" "crypto/cipher" "crypto/sha256" "crypto/tls" "encoding/binary" "encoding/hex" "encoding/json" "errors" "fmt" "io" "math" "net" "net/http" "net/http/httptest" "os" "path/filepath" "reflect" "regexp" "runtime" "runtime/debug" runtimemetrics "runtime/metrics" "sort" "strconv" "strings" "sync" "sync/atomic" "syscall" "testing" "time" protocol "git.sechmachine.io.vn/sechmachine/VerseVDI-Protocol/gen/go/protocol" ) const ( qualificationToolVersion = "versevdi-gateway-qualification/v6" qualificationImpairmentQueuePackets = nativeApolloVideoQueuePackets qualificationImpairmentMaxPackets = 100_000 qualificationImpairmentPacketCount = 10_000 qualificationProcessingLimit = 5 * time.Millisecond qualificationImpairmentSeed uint64 = 0x3c6a11ce qualificationClockOverheadMethod = "median of 1000 batches of 100 monotonic time reads" qualificationGatewayCPUScope = "isolated gateway subprocess; bounded recorder/control included, fixture and client driver excluded" qualificationResourceMethod = "RUSAGE_SELF user+system CPU; runtime/metrics heap objects, allocated objects/bytes, and live goroutines sampled once per second" ) type qualificationMediaProfile struct { Name string Codec string BitrateKbps int64 FPS int Duration time.Duration Warmup time.Duration PacketBytes int } type qualificationPathTrace struct { NativeSetup bool NativeOpen bool NativeUDPIngress bool ApolloRecovered bool ProductionQueue bool ProductionMediaLoop bool ProductionPacer bool VerseQUIC bool PublicClientDecode bool PayloadPreserved bool } type qualificationPath struct { client *independentGatewayClient server *Server session *nativeApolloSession fixture *qualificationApolloFixture backend *qualificationTracingBackend process *qualificationGatewayProcess key []byte flow string frame uint32 bootTrace atomic.Bool sourceUDP atomic.Uint64 closeOnce sync.Once shutdown func() } type qualificationImpairmentProfile struct { Name string RTT time.Duration Jitter time.Duration LossPercent float64 Reorder bool CapacitySteps []int } type qualificationProcessingSummary struct { Profile string `json:"profile"` Codec string `json:"codec"` ConfiguredBitrateKbps int64 `json:"configured_bitrate_kbps"` ObservedBitrateKbps float64 `json:"observed_bitrate_kbps"` ConfiguredFPS int `json:"configured_fps"` ObservedFPS float64 `json:"observed_fps"` PayloadBytes int64 `json:"payload_bytes"` MinimumFrameBytes int `json:"minimum_frame_bytes"` MaximumFrameBytes int `json:"maximum_frame_bytes"` Warmup time.Duration `json:"warmup_ns"` ConfiguredDuration time.Duration `json:"configured_duration_ns"` ActualDuration time.Duration `json:"actual_duration_ns"` Count int64 `json:"count"` Min time.Duration `json:"min_ns"` Median time.Duration `json:"median_ns"` P90 time.Duration `json:"p90_ns"` P95 time.Duration `json:"p95_ns"` P99 time.Duration `json:"p99_ns"` Max time.Duration `json:"max_ns"` Mean time.Duration `json:"mean_ns"` StandardDeviation time.Duration `json:"standard_deviation_ns"` ClockOverhead time.Duration `json:"clock_overhead_ns"` ClockMethod string `json:"clock_overhead_method"` Histogram map[string]int `json:"histogram"` PayloadSHA256 string `json:"payload_sha256"` RawSamples string `json:"raw_samples"` RawSamplesSHA256 string `json:"raw_samples_sha256"` RawSamplesBytes int64 `json:"raw_samples_bytes"` ResourceSamples int `json:"resource_samples"` CPUSeconds float64 `json:"cpu_seconds"` CPUScope string `json:"cpu_scope"` ResourceMethod string `json:"resource_method"` PeakHeapBytes uint64 `json:"peak_heap_bytes"` PeakGoroutines int `json:"peak_goroutines"` AllocatedObjects uint64 `json:"allocated_objects"` AllocatedBytes uint64 `json:"allocated_bytes"` RawResources string `json:"raw_resources"` RawResourcesSHA256 string `json:"raw_resources_sha256"` RawResourcesBytes int64 `json:"raw_resources_bytes"` } type qualificationImpairmentObservation struct { Profile string `json:"profile"` MediaProfile string `json:"media_profile"` Seed uint64 `json:"seed"` Sent int `json:"sent"` Delivered int `json:"delivered"` Dropped int `json:"dropped"` InjectedDropped int `json:"injected_dropped"` InjectedReordered int `json:"injected_reordered"` SourceEmitted int `json:"source_emitted"` SourceDatagrams uint64 `json:"source_datagrams"` ProviderIngressDatagrams uint64 `json:"provider_ingress_datagrams"` ProviderRecovered int `json:"provider_recovered"` ProviderFECDropped int `json:"provider_fec_dropped"` ProviderEnqueued int `json:"provider_enqueued"` ProviderEnqueueDropped int `json:"provider_enqueue_dropped"` ProviderQueueReplaced int `json:"provider_queue_replaced"` GatewayForwarded int `json:"gateway_forwarded"` GatewayDropped int `json:"gateway_dropped"` QUICSent int `json:"quic_sent"` QUICDatagramsSent int `json:"quic_datagrams_sent"` QUICSendDropped int `json:"quic_send_dropped"` ClientDeliveryDropped int `json:"client_delivery_dropped"` UnexplainedDropped int `json:"unexplained_dropped"` ObservedOutOfOrder int `json:"observed_out_of_order"` ObservedLatency time.Duration `json:"observed_one_way_latency_ns"` ObservedRTT time.Duration `json:"observed_rtt_ns"` RTTSource string `json:"rtt_source"` ObservedJitter time.Duration `json:"observed_jitter_ns"` ObservedLossPercent float64 `json:"observed_loss_percent"` ObservedReorderPercent float64 `json:"observed_reorder_percent"` ObservedThroughputKbps float64 `json:"observed_throughput_kbps"` MaxQueuePackets int `json:"max_queue_packets"` ConfiguredRTT time.Duration `json:"configured_rtt_ns"` ConfiguredJitter time.Duration `json:"configured_jitter_ns"` AppliedJitter time.Duration `json:"applied_jitter_ns"` ConfiguredLossPercent float64 `json:"configured_loss_percent"` ConfiguredReorder bool `json:"configured_reorder"` ConfiguredCapacitySteps []int `json:"configured_capacity_steps_percent"` CapacityStepObservations []qualificationCapacityStep `json:"capacity_step_observations,omitempty"` RawSamples string `json:"raw_samples"` RawSamplesSHA256 string `json:"raw_samples_sha256"` RawSamplesBytes int64 `json:"raw_samples_bytes"` } type qualificationCapacityStep struct { ReductionPercent int `json:"reduction_percent"` Convergence time.Duration `json:"convergence_ns"` MaximumFiveSecond int64 `json:"maximum_five_second_bytes"` FiveSecondCap int64 `json:"five_second_cap_bytes"` } type qualificationResourceSample struct { Elapsed time.Duration CPUSeconds float64 HeapBytes uint64 Goroutines int AllocatedObjects uint64 AllocatedBytes uint64 } type qualificationDeliverySample struct { At time.Time Bytes int64 } type qualificationFairnessEvidence struct { Evaluation time.Duration `json:"evaluation_ns"` PerFlowBytes map[string]int64 `json:"per_flow_bytes"` ShareError map[string]float64 `json:"share_error"` JainIndex float64 `json:"jain_index"` CapacitySteps []qualificationCapacityStep `json:"capacity_steps"` Series []qualificationFairnessSeries `json:"per_flow_aggregate_series"` RawSamples string `json:"raw_samples"` RawSamplesSHA256 string `json:"raw_samples_sha256"` RawSamplesBytes int64 `json:"raw_samples_bytes"` } type qualificationFairnessSeries struct { Elapsed time.Duration `json:"elapsed_ns"` PerFlowBytes map[string]int64 `json:"per_flow_bytes"` AggregateBytes int64 `json:"aggregate_bytes"` } type qualificationManifest struct { Status string `json:"status"` ToolVersion string `json:"tool_version"` Command string `json:"command"` CandidateCommit string `json:"candidate_commit"` ProtocolVersion string `json:"protocol_version"` StartedAt string `json:"started_at"` CompletedAt string `json:"completed_at"` GoVersion string `json:"go_version"` ToolVersions map[string]string `json:"tool_versions"` OS string `json:"os"` Architecture string `json:"architecture"` Topology string `json:"topology"` Direction string `json:"direction"` QueueDiscipline string `json:"queue_discipline"` Evidence []string `json:"evidence_classification"` Deferred []string `json:"deferred"` Processing []qualificationProcessingSummary `json:"processing"` Impairments []qualificationImpairmentObservation `json:"impairments"` Fairness qualificationFairnessEvidence `json:"fairness"` } func qualificationMediaProfiles() []qualificationMediaProfile { return []qualificationMediaProfile{ {Name: "1080p60-h264", Codec: "h264", BitrateKbps: 20000, FPS: 60, Duration: 10 * time.Minute, Warmup: time.Second, PacketBytes: 1179}, {Name: "1440p120-hevc", Codec: "hevc", BitrateKbps: 50000, FPS: 120, Duration: 10 * time.Minute, Warmup: time.Second, PacketBytes: 1179}, {Name: "4k60-hevc", Codec: "hevc", BitrateKbps: 80000, FPS: 60, Duration: 10 * time.Minute, Warmup: time.Second, PacketBytes: 1179}, } } func qualificationImpairmentProfiles() []qualificationImpairmentProfile { return []qualificationImpairmentProfile{ {Name: "baseline", RTT: 20 * time.Millisecond}, {Name: "latency", RTT: 150 * time.Millisecond}, {Name: "jitter", RTT: 50 * time.Millisecond, Jitter: 30 * time.Millisecond}, {Name: "loss", RTT: 50 * time.Millisecond, Jitter: 10 * time.Millisecond, LossPercent: 5}, {Name: "reorder", RTT: 100 * time.Millisecond, Jitter: 10 * time.Millisecond, LossPercent: 1, Reorder: true}, {Name: "constrained", RTT: 50 * time.Millisecond, Jitter: 10 * time.Millisecond, LossPercent: 2, Reorder: true, CapacitySteps: []int{25, 50}}, } } func validateQualificationOutputDir(path string) error { if path == "" || !filepath.IsAbs(path) || filepath.Clean(path) == string(filepath.Separator) { return errors.New("qualification evidence directory must be a non-root absolute path") } return nil } func qualificationPayload(profile qualificationMediaProfile) []byte { payload := make([]byte, profile.PacketBytes) if profile.Codec == "h264" { copy(payload, []byte{0, 0, 1, 0x65}) } else { copy(payload, []byte{0, 0, 1, 0x26}) } for index := 4; index < len(payload); index++ { payload[index] = byte(index*31 + len(profile.Name)) } return payload } func qualificationFrameRate(profile qualificationMediaProfile) int { if profile.FPS > 0 { return profile.FPS } return 60 } func qualificationFrameSize(profile qualificationMediaProfile, index int64) int { fps := qualificationFrameRate(profile) bytesPerSecond := int(profile.BitrateKbps * 1000 / 8) position := int(index % int64(fps)) if fps == 1 { return min(bytesPerSecond, maxCompleteFrameBytes) } keyframeBytes := min(bytesPerSecond/fps*4, maxCompleteFrameBytes) if position == 0 { return keyframeBytes } remaining := bytesPerSecond - keyframeBytes size := remaining / (fps - 1) if position <= remaining%(fps-1) { size++ } return size } func qualificationFramePayload(profile qualificationMediaProfile, index int64) []byte { payload := make([]byte, qualificationFrameSize(profile, index)) if profile.Codec == "h264" { copy(payload, []byte{0, 0, 1, 0x65}) } else { copy(payload, []byte{0, 0, 1, 0x26}) } for offset := 4; offset < len(payload); offset++ { payload[offset] = byte(int64(offset)*31 + index*17 + int64(len(profile.Name))) } if len(payload) >= 8 { binary.BigEndian.PutUint64(payload[len(payload)-8:], uint64(index)) } return payload } func qualificationMediaPacerKbps(profile qualificationMediaProfile, reduction int) int64 { payloadKbps := profile.BitrateKbps * int64(100-reduction) / 100 datagrams := (profile.PacketBytes + frameV2PayloadSize - 1) / frameV2PayloadSize wireBytes := profile.PacketBytes + datagrams*frameV2HeaderSize return (payloadKbps*int64(wireBytes) + int64(profile.PacketBytes) - 1) / int64(profile.PacketBytes) } func qualificationFramePacerKbps(profile qualificationMediaProfile) int64 { fps := qualificationFrameRate(profile) var payloadBytes, wireBytes int64 for index := range fps { size := qualificationFrameSize(profile, int64(index)) datagrams := (size + frameV2PayloadSize - 1) / frameV2PayloadSize payloadBytes += int64(size) wireBytes += int64(size + datagrams*frameV2HeaderSize) } return (profile.BitrateKbps*wireBytes + payloadBytes - 1) / payloadBytes } func qualificationBoundedRelease(target, next, now time.Time, spacing time.Duration) time.Time { release := target if next.After(release) { release = next } if now.Sub(release) > spacing { return now.Add(-spacing) } return release } func qualificationWaitUntil(target time.Time) { const preciseWindow = 500 * time.Microsecond for { delay := time.Until(target) if delay <= 0 { return } if delay > preciseWindow { time.Sleep(delay - preciseWindow) } else if delay > 50*time.Microsecond { runtime.Gosched() } } } type qualificationTracingBackend struct { native *NativeApolloBackend setups atomic.Uint64 opens atomic.Uint64 sessions sync.Map } func (b *qualificationTracingBackend) Management(ctx context.Context, request LaunchRequest) ([]byte, error) { return b.native.Management(ctx, request) } func (b *qualificationTracingBackend) Setup(ctx context.Context, request LaunchRequest, management []byte) ([]byte, error) { response, err := b.native.Setup(ctx, request, management) if err == nil { b.setups.Add(1) } return response, err } func (b *qualificationTracingBackend) Open(ctx context.Context, request LaunchRequest, response RTSPResponse) (ProviderSession, error) { session, err := b.native.Open(ctx, request, response) if err == nil { native, ok := session.(*nativeApolloSession) if !ok { return nil, ErrProviderMalformed } b.opens.Add(1) b.sessions.Store(request.SessionID, native) } return session, err } func (b *qualificationTracingBackend) session(sessionID string) *nativeApolloSession { value, _ := b.sessions.Load(sessionID) session, _ := value.(*nativeApolloSession) return session } type qualificationApolloFixture struct { sessionID string management *httptest.Server stream net.Listener control *net.UDPConn audio *net.UDPConn video *net.UDPConn videoRemote atomic.Pointer[net.UDPAddr] key atomic.Pointer[[]byte] keyReady chan []byte failures chan error closed atomic.Bool closeOnce sync.Once sentPackets atomic.Uint64 work protocol.ProviderSessionWork controlImpairmentMu sync.Mutex controlRTT time.Duration controlJitter time.Duration controlRandom uint64 } func newQualificationApolloFixture(t *testing.T, serverTLS, clientTLS *tls.Config, sessionID string, profile qualificationMediaProfile) *qualificationApolloFixture { t.Helper() fixture := &qualificationApolloFixture{sessionID: sessionID, keyReady: make(chan []byte, 1), failures: make(chan error, 8)} var err error fixture.stream, err = net.Listen("tcp", "127.0.0.1:0") if err != nil { t.Fatal(err) } fixture.control, err = net.ListenUDP("udp", &net.UDPAddr{IP: net.ParseIP("127.0.0.1")}) if err != nil { fixture.Close() t.Fatal(err) } fixture.audio, err = net.ListenUDP("udp", &net.UDPAddr{IP: net.ParseIP("127.0.0.1")}) if err != nil { fixture.Close() t.Fatal(err) } fixture.video, err = net.ListenUDP("udp", &net.UDPAddr{IP: net.ParseIP("127.0.0.1")}) if err != nil { fixture.Close() t.Fatal(err) } streamHost, streamPortText, err := net.SplitHostPort(fixture.stream.Addr().String()) if err != nil { fixture.Close() t.Fatal(err) } streamPort, err := strconv.ParseInt(streamPortText, 10, 64) if err != nil { fixture.Close() t.Fatal(err) } fixture.management = httptest.NewUnstartedServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { if request.TLS == nil || len(request.TLS.PeerCertificates) != 1 { http.Error(response, "mTLS required", http.StatusUnauthorized) return } switch request.URL.Path { case "/serverinfo": _, _ = response.Write([]byte("qualification-apollo2571869449984")) case "/applist": _, _ = response.Write([]byte("42")) case "/launch": key, decodeErr := hex.DecodeString(request.URL.Query().Get("rikey")) if decodeErr != nil || len(key) != 16 { http.Error(response, "invalid RIK", http.StatusBadRequest) return } copyKey := append([]byte(nil), key...) fixture.key.Store(©Key) select { case fixture.keyReady <- copyKey: default: } _, _ = response.Write([]byte("rtspenc://" + fixture.stream.Addr().String() + "")) default: http.NotFound(response, request) } })) fixture.management.TLS = serverTLS fixture.management.StartTLS() managementHost, managementPortText, err := net.SplitHostPort(fixture.management.Listener.Addr().String()) if err != nil { fixture.Close() t.Fatal(err) } managementPort, err := strconv.ParseInt(managementPortText, 10, 64) if err != nil { fixture.Close() t.Fatal(err) } codec, width, height, fps := "H264", int64(1920), int64(1080), int64(60) if profile.Codec == "hevc" { codec = "HEVC" if strings.Contains(profile.Name, "1440") { width, height, fps = 2560, 1440, 120 } else { width, height, fps = 3840, 2160, 60 } } pinned := sha256.Sum256(serverTLS.Certificates[0].Certificate[0]) fixture.work = protocol.ProviderSessionWork{ Version: "1", SessionID: sessionID, GatewayID: "gateway-1", ExpiresAt: "2099-01-01T00:00:00Z", ProviderProfile: ProviderProfileApollo, ProviderIdentity: "qualification-apollo#sha256:" + hex.EncodeToString(pinned[:]), PolicyVersionID: "qualification-policy", ApplicationID: "42", ClientID: "qualification-client", StreamPolicy: protocol.ProviderStreamPolicy{ ResolutionWidth: width, ResolutionHeight: height, Fps: fps, Codec: codec, BitrateKbps: profile.BitrateKbps, AudioEnabled: true, }, ManagementHost: managementHost, ManagementPort: managementPort, StreamHost: streamHost, StreamPort: streamPort, ClientCertificatePem: certificatePEM(t, clientTLS.Certificates[0]), ClientPrivateKeyPem: privateKeyPEM(t, clientTLS.Certificates[0]), ServerCertificatePem: certificatePEM(t, tls.Certificate{Certificate: [][]byte{serverTLS.Certificates[0].Certificate[1]}}), ClipboardPolicy: protocol.ClipboardPolicy{MaxTextBytes: 65536, MaxUpdatesPerMinute: 30}, } go fixture.serveRTSP() go fixture.serveControl() go fixture.serveMedia(fixture.audio, false) go fixture.serveMedia(fixture.video, true) t.Cleanup(fixture.Close) return fixture } func (f *qualificationApolloFixture) fail(err error) { if err == nil || f.closed.Load() { return } select { case f.failures <- err: default: } } func (f *qualificationApolloFixture) serveRTSP() { key := <-f.keyReady codec, err := newEncryptedRTSPCodec(key) if err != nil { f.fail(err) return } methods := []string{"OPTIONS", "DESCRIBE", "SETUP", "SETUP", "SETUP", "ANNOUNCE", "PLAY"} targets := []string{"rtspenc://" + f.stream.Addr().String(), "rtspenc://" + f.stream.Addr().String(), "streamid=audio/0/0", "streamid=video/0/0", "streamid=control/13/0", "streamid=control/13/0", "/"} describe := "a=x-ss-general.featureFlags:1\r\na=x-ss-general.encryptionSupported:7\r\na=x-ss-general.encryptionRequested:1\r\na=fmtp:97 surround-params=21101\r\n" responses := []string{ "RTSP/1.0 200 OK\r\nCSeq: 1\r\n\r\n", fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 2\r\nContent-Type: application/sdp\r\nContent-Length: %d\r\n\r\n%s", len(describe), describe), fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 3\r\nSession: %s;timeout=90\r\nTransport: unicast;server_port=%d\r\nX-SS-Ping-Payload: 0123456789abcdef\r\n\r\n", f.sessionID, f.audio.LocalAddr().(*net.UDPAddr).Port), fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 4\r\nSession: %s\r\nTransport: unicast;server_port=%d\r\nX-SS-Ping-Payload: fedcba9876543210\r\n\r\n", f.sessionID, f.video.LocalAddr().(*net.UDPAddr).Port), fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 5\r\nSession: %s\r\nTransport: unicast;server_port=%d\r\nX-SS-Connect-Data: 305419896\r\n\r\n", f.sessionID, f.control.LocalAddr().(*net.UDPAddr).Port), fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 6\r\nSession: %s\r\n\r\n", f.sessionID), fmt.Sprintf("RTSP/1.0 200 OK\r\nCSeq: 7\r\nSession: %s\r\n\r\n", f.sessionID), } for index, method := range methods { connection, acceptErr := f.stream.Accept() if acceptErr != nil { f.fail(acceptErr) return } header := make([]byte, encryptedRTSPHeaderSize) if _, err = io.ReadFull(connection, header); err != nil { _ = connection.Close() f.fail(err) return } length := binary.BigEndian.Uint32(header[:4]) & 0x7fffffff if length > maxApolloRTSPHeaders+maxApolloRTSPBody { _ = connection.Close() f.fail(ErrProviderMalformed) return } frame := make([]byte, encryptedRTSPHeaderSize+int(length)) copy(frame, header) if _, err = io.ReadFull(connection, frame[encryptedRTSPHeaderSize:]); err != nil { _ = connection.Close() f.fail(err) return } plaintext, openErr := codec.OpenClient(frame) firstLine, _, _ := strings.Cut(string(plaintext), "\r\n") if openErr != nil || firstLine != method+" "+targets[index]+" RTSP/1.0" { _ = connection.Close() f.fail(ErrProviderMalformed) return } responseFrame, sealErr := hostEncryptedRTSPFrameNoTest(key, uint32(index+1), []byte(responses[index])) if sealErr != nil { _ = connection.Close() f.fail(sealErr) return } _, err = connection.Write(responseFrame) _ = connection.Close() if err != nil { f.fail(err) return } } } func hostEncryptedRTSPFrameNoTest(key []byte, sequence uint32, plaintext []byte) ([]byte, error) { block, err := aes.NewCipher(key) if err != nil { return nil, err } aead, err := cipher.NewGCM(block) if err != nil { return nil, err } nonce := encryptedRTSPNonce(sequence, 'H', 'R') sealed := aead.Seal(nil, nonce[:], plaintext, nil) frame := make([]byte, encryptedRTSPHeaderSize+len(plaintext)) binary.BigEndian.PutUint32(frame[:4], uint32(len(plaintext))|0x80000000) binary.BigEndian.PutUint32(frame[4:8], sequence) copy(frame[8:24], sealed[len(plaintext):]) copy(frame[24:], sealed[:len(plaintext)]) return frame, nil } func (f *qualificationApolloFixture) serveControl() { buffer := make([]byte, apolloENetMaximumPacket) count, remote, err := f.control.ReadFromUDP(buffer) if err != nil { f.fail(err) return } connect := buffer[:count] if count != 52 || connect[4] != apolloENetConnect|apolloENetAcknowledged || binary.BigEndian.Uint32(connect[20:24]) != apolloENetChannels { f.fail(ErrProviderMalformed) return } verify := make([]byte, 48) apolloENetHeader(verify, 0, 0, time.Now()) verify[4] = apolloENetVerifyConnect | apolloENetAcknowledged verify[5] = 0xff binary.BigEndian.PutUint16(verify[6:8], 1) binary.BigEndian.PutUint16(verify[8:10], 7) verify[10], verify[11] = 2, 3 binary.BigEndian.PutUint32(verify[12:16], 1400) binary.BigEndian.PutUint32(verify[16:20], 32768) binary.BigEndian.PutUint32(verify[20:24], apolloENetChannels) binary.BigEndian.PutUint32(verify[44:48], binary.BigEndian.Uint32(connect[44:48])) if _, err = f.control.WriteToUDP(verify, remote); err != nil { f.fail(err) return } for { count, _, err = f.control.ReadFromUDP(buffer) if err != nil { f.fail(err) return } packet := buffer[:count] if len(packet) < 8 { f.fail(ErrProviderMalformed) return } command, channel := packet[4]&apolloENetCommandMask, packet[5] sequence := binary.BigEndian.Uint16(packet[6:8]) switch command { case 1: case apolloENetSendReliable, apolloENetPing, apolloENetDisconnect: f.sendControlAcknowledge(sourceShapedENetAcknowledgePacket(7, 2, channel, sequence), remote) if command == apolloENetDisconnect { return } case apolloENetSendUnsequenced: default: f.fail(ErrProviderMalformed) return } } } func (f *qualificationApolloFixture) setControlImpairment(profile qualificationImpairmentProfile) { f.controlImpairmentMu.Lock() f.controlRTT = profile.RTT f.controlJitter = profile.Jitter f.controlRandom = qualificationImpairmentSeed f.controlImpairmentMu.Unlock() } func (f *qualificationApolloFixture) controlResponseDelay() time.Duration { f.controlImpairmentMu.Lock() defer f.controlImpairmentMu.Unlock() delay := f.controlRTT if f.controlJitter > 0 { f.controlRandom ^= f.controlRandom << 13 f.controlRandom ^= f.controlRandom >> 7 f.controlRandom ^= f.controlRandom << 17 width := uint64(f.controlJitter*2 + 1) delay += time.Duration(f.controlRandom%width) - f.controlJitter } return max(delay, 0) } func (f *qualificationApolloFixture) sendControlAcknowledge(packet []byte, remote *net.UDPAddr) { delay := f.controlResponseDelay() copyPacket := append([]byte(nil), packet...) copyRemote := *remote go func() { timer := time.NewTimer(delay) defer timer.Stop() <-timer.C if f.closed.Load() { return } if _, err := f.control.WriteToUDP(copyPacket, ©Remote); err != nil { f.fail(err) } }() } func (f *qualificationApolloFixture) serveMedia(socket *net.UDPConn, video bool) { buffer := make([]byte, apolloMediaMaximumPacket) for { _, remote, err := socket.ReadFromUDP(buffer) if err != nil { f.fail(err) return } if video && f.videoRemote.Load() == nil { copyRemote := *remote f.videoRemote.Store(©Remote) } } } func (f *qualificationApolloFixture) sendVideo(ctx context.Context, packets [][]byte) error { for f.videoRemote.Load() == nil { select { case err := <-f.failures: return err case <-ctx.Done(): return ctx.Err() case <-time.After(time.Millisecond): } } remote := f.videoRemote.Load() for _, packet := range packets { if _, err := f.video.WriteToUDP(packet, remote); err != nil { return err } f.sentPackets.Add(1) } return nil } func (f *qualificationApolloFixture) streamKey() ([]byte, error) { key := f.key.Load() if key == nil || len(*key) != 16 { return nil, ErrProviderMalformed } return append([]byte(nil), (*key)...), nil } func (f *qualificationApolloFixture) Close() { if f == nil { return } f.closeOnce.Do(func() { f.closed.Store(true) if f.management != nil { f.management.Close() } if f.stream != nil { _ = f.stream.Close() } for _, socket := range []*net.UDPConn{f.control, f.audio, f.video} { if socket != nil { _ = socket.Close() } } }) } type qualificationAdmissionRecord struct { authority protocol.SessionAuthority work protocol.ProviderSessionWork used bool } type qualificationAdmission struct { mu sync.Mutex records map[string]*qualificationAdmissionRecord } func (a *qualificationAdmission) Admit(_ context.Context, request protocol.TunnelAdmissionRequest) (protocol.SessionAuthority, error) { a.mu.Lock() defer a.mu.Unlock() record := a.records[request.SessionID] if record == nil || record.used || request.GatewayID != record.authority.GatewayID || request.Audience != record.authority.Audience || !reflect.DeepEqual(request.Capabilities, record.authority.Capabilities) { return protocol.SessionAuthority{}, ErrAdmissionRejected } record.used = true return record.authority, nil } func (a *qualificationAdmission) ProviderWork(_ context.Context, authority protocol.SessionAuthority) (protocol.ProviderSessionWork, error) { a.mu.Lock() defer a.mu.Unlock() record := a.records[authority.SessionID] if record == nil || !record.used || !reflect.DeepEqual(authority, record.authority) { return protocol.ProviderSessionWork{}, ErrAdmissionRejected } return record.work, nil } func (*qualificationAdmission) Release(context.Context, protocol.SessionAuthority) error { return nil } type qualificationFleet struct { t *testing.T server *Server paths []*qualificationPath clients []*independentGatewayClient cancel context.CancelFunc serveDone chan error closeOnce sync.Once } func newQualificationFleet(t *testing.T, count int, profile qualificationMediaProfile, pacerKbps int64) *qualificationFleet { t.Helper() serverTLS, clientTLS := testTLS(t) backend := &qualificationTracingBackend{native: NewNativeApolloBackend()} admission := &qualificationAdmission{records: make(map[string]*qualificationAdmissionRecord, count)} fixtures := make([]*qualificationApolloFixture, 0, count) for index := 0; index < count; index++ { sessionID := fmt.Sprintf("qualification-flow-%d", index+1) fixture := newQualificationApolloFixture(t, serverTLS, clientTLS, sessionID, profile) authority := protocol.SessionAuthority{ Version: "1", SessionID: sessionID, GatewayID: "gateway-1", Audience: "versevdi-gateway", ExpiresAt: time.Now().Add(2 * time.Minute).UTC().Format(time.RFC3339Nano), Capabilities: DefaultCapabilities(), ProviderProfile: ProviderProfileApollo, ProviderIdentity: fixture.work.ProviderIdentity, } work := fixture.work work.ExpiresAt = authority.ExpiresAt admission.records[sessionID] = &qualificationAdmissionRecord{authority: authority, work: work} fixtures = append(fixtures, fixture) } server, err := NewServer(ServerConfig{ ListenAddress: "127.0.0.1:0", TLSConfig: serverTLS, GatewayID: "gateway-1", Capabilities: DefaultCapabilities(), ProviderCapabilities: DefaultCapabilities(), Admission: admission, ProviderStateReporter: &recordingProviderStateReporter{}, Provider: NewApolloAdapter(backend, ProviderIdentity{}), PacerKbps: pacerKbps, }) if err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) fleet := &qualificationFleet{t: t, server: server, cancel: cancel, serveDone: make(chan error, 1)} go func() { fleet.serveDone <- server.Serve(ctx) }() for index, fixture := range fixtures { sessionID := fmt.Sprintf("qualification-flow-%d", index+1) request := protocol.TunnelAdmissionRequest{ Version: "1", SessionID: sessionID, GatewayID: "gateway-1", Audience: "versevdi-gateway", Grant: strings.Repeat(string(rune('a'+index)), 64), ClientNonce: fmt.Sprintf("nonce-fleet-%06d", index), DeviceSignature: strings.Repeat("s", 86), Capabilities: DefaultCapabilities(), } client, err := dialIndependentGateway(context.Background(), server.Addr().String(), clientTLS, request) if err != nil { fleet.Close() t.Fatal(err) } session := backend.session(sessionID) key, keyErr := fixture.streamKey() if keyErr != nil || session == nil { fleet.Close() t.Fatalf("qualification fleet session %s unavailable: %v", sessionID, keyErr) } fleet.clients = append(fleet.clients, client) fleet.paths = append(fleet.paths, &qualificationPath{ client: client, server: server, session: session, fixture: fixture, backend: backend, key: key, shutdown: func() {}, }) } t.Cleanup(fleet.Close) return fleet } func (f *qualificationFleet) Close() { if f == nil { return } f.closeOnce.Do(func() { for _, client := range f.clients { _ = client.Close() } f.cancel() _ = f.server.Close() if err := <-f.serveDone; err != nil { f.t.Errorf("serve qualification fleet: %v", err) } }) } func newQualificationPath(t *testing.T, profile qualificationMediaProfile, pacerKbps int64) *qualificationPath { return newQualificationPathWithImpairment(t, profile, pacerKbps, nil) } func newQualificationImpairedPath(t *testing.T, profile qualificationMediaProfile, pacerKbps int64, impairment qualificationImpairmentProfile) *qualificationPath { return newQualificationPathWithImpairment(t, profile, pacerKbps, &impairment) } func newQualificationPathWithImpairment(t *testing.T, profile qualificationMediaProfile, pacerKbps int64, impairment *qualificationImpairmentProfile) *qualificationPath { t.Helper() serverTLS, clientTLS := testTLS(t) fixture := newQualificationApolloFixture(t, serverTLS, clientTLS, "qualification-session", profile) if impairment != nil { fixture.setControlImpairment(*impairment) } backend := &qualificationTracingBackend{native: NewNativeApolloBackend()} provider := NewApolloAdapter(backend, ProviderIdentity{}) authority := protocol.SessionAuthority{ Version: "1", SessionID: "qualification-session", GatewayID: "gateway-1", Audience: "versevdi-gateway", ReconnectSequence: 0, ExpiresAt: time.Now().Add(45 * time.Minute).UTC().Format(time.RFC3339Nano), Capabilities: DefaultCapabilities(), ProviderProfile: ProviderProfileApollo, ProviderIdentity: fixture.work.ProviderIdentity, } work := fixture.work work.ExpiresAt = authority.ExpiresAt admission := &oneTimeAdmission{authority: authority, released: make(chan struct{}), disableClipboard: true, providerWork: &work} server, err := NewServer(ServerConfig{ ListenAddress: "127.0.0.1:0", TLSConfig: serverTLS, GatewayID: authority.GatewayID, Capabilities: DefaultCapabilities(), ProviderCapabilities: DefaultCapabilities(), Admission: admission, ProviderStateReporter: &recordingProviderStateReporter{}, Provider: provider, PacerKbps: pacerKbps, }) if err != nil { t.Fatal(err) } ctx, cancel := context.WithCancel(context.Background()) serveDone := make(chan error, 1) go func() { serveDone <- server.Serve(ctx) }() request := protocol.TunnelAdmissionRequest{ Version: "1", SessionID: authority.SessionID, GatewayID: authority.GatewayID, Audience: authority.Audience, Grant: strings.Repeat("g", 64), ClientNonce: "nonce-qualification", DeviceSignature: strings.Repeat("s", 86), Capabilities: DefaultCapabilities(), } client, err := dialIndependentGateway(context.Background(), server.Addr().String(), clientTLS, request) if err != nil { cancel() _ = server.Close() t.Fatal(err) } session := backend.session(authority.SessionID) key, err := fixture.streamKey() if err != nil || session == nil { _ = client.Close() cancel() _ = server.Close() t.Fatalf("native qualification session unavailable: %v", err) } path := &qualificationPath{client: client, server: server, session: session, fixture: fixture, backend: backend, key: key} path.shutdown = func() { _ = client.Close() cancel() _ = server.Close() if err := <-serveDone; err != nil { t.Errorf("serve qualification path: %v", err) } } t.Cleanup(path.Close) return path } func newQualificationProcessingPath(t *testing.T, profile qualificationMediaProfile, pacerKbps int64) *qualificationPath { t.Helper() serverTLS, clientTLS := testTLS(t) fixture := newQualificationApolloFixture(t, serverTLS, clientTLS, "qualification-session", profile) authority := protocol.SessionAuthority{ Version: "1", SessionID: "qualification-session", GatewayID: "gateway-1", Audience: "versevdi-gateway", ReconnectSequence: 0, ExpiresAt: time.Now().Add(45 * time.Minute).UTC().Format(time.RFC3339Nano), Capabilities: DefaultCapabilities(), ProviderProfile: ProviderProfileApollo, ProviderIdentity: fixture.work.ProviderIdentity, } work := fixture.work work.ExpiresAt = authority.ExpiresAt process := startQualificationGatewayProcess(t, serverTLS, authority, work, pacerKbps) request := protocol.TunnelAdmissionRequest{ Version: "1", SessionID: authority.SessionID, GatewayID: authority.GatewayID, Audience: authority.Audience, Grant: strings.Repeat("g", 64), ClientNonce: "nonce-qualification", DeviceSignature: strings.Repeat("s", 86), Capabilities: DefaultCapabilities(), } client, err := dialIndependentGateway(context.Background(), process.ready.GatewayAddress, clientTLS, request) if err != nil { process.Close() t.Fatal(err) } key, err := fixture.streamKey() if err != nil { _ = client.Close() process.Close() t.Fatal(err) } path := &qualificationPath{client: client, fixture: fixture, process: process, key: key} path.shutdown = func() { _ = client.Close() process.Close() } t.Cleanup(path.Close) return path } func (p *qualificationPath) Close() { if p != nil { p.closeOnce.Do(p.shutdown) } } func (p *qualificationPath) traverse(t *testing.T, payload []byte) (qualificationPathTrace, time.Duration, error) { t.Helper() beforeMetrics := p.server.Metrics() beforeIngress := p.session.mediaIngress.Load() beforeRecovered := p.session.mediaRecovered.Load() beforeEnqueued := p.session.mediaEnqueued.Load() beforePacer := p.server.pacer.reservations.Load() trace, err := p.emit(t, payload) if err != nil { return trace, 0, err } recovered, err := p.receivePayload(context.Background()) if err != nil { return trace, 0, err } afterMetrics := p.server.Metrics() deadline := time.Now().Add(2 * time.Second) for (afterMetrics.ProcessingSamples < beforeMetrics.ProcessingSamples+1 || afterMetrics.MediaPackets <= beforeMetrics.MediaPackets || p.session.mediaEnqueued.Load() <= beforeEnqueued) && time.Now().Before(deadline) { runtime.Gosched() afterMetrics = p.server.Metrics() } trace.NativeUDPIngress = p.session.mediaIngress.Load() > beforeIngress trace.ApolloRecovered = p.session.mediaRecovered.Load() > beforeRecovered trace.ProductionQueue = p.session.mediaEnqueued.Load() > beforeEnqueued trace.ProductionMediaLoop = afterMetrics.ProcessingSamples > beforeMetrics.ProcessingSamples trace.ProductionPacer = p.server.pacer.reservations.Load() > beforePacer trace.VerseQUIC = afterMetrics.MediaPackets > beforeMetrics.MediaPackets trace.PublicClientDecode = true trace.PayloadPreserved = bytes.Equal(recovered, payload) if p.bootTrace.CompareAndSwap(false, true) { trace.NativeSetup = p.backend.setups.Load() == 1 trace.NativeOpen = p.backend.opens.Load() == 1 } return trace, time.Duration(afterMetrics.ProcessingDelayNanos - beforeMetrics.ProcessingDelayNanos), nil } func (p *qualificationPath) emit(t *testing.T, payload []byte) (qualificationPathTrace, error) { t.Helper() if p == nil || p.fixture == nil || len(payload) == 0 || len(payload) > apolloVideoMaximumBlocks*apolloVideoMaximumDataShards*apolloVideoShardPayloadSize-8 { return qualificationPathTrace{}, ErrProviderMalformed } p.frame++ packets := qualificationSourceVideoPackets(t, p.key, p.frame, payload) ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() if err := p.fixture.sendVideo(ctx, packets); err != nil { return qualificationPathTrace{}, err } p.sourceUDP.Add(uint64(len(packets))) return qualificationPathTrace{}, nil } func (p *qualificationPath) receivePayload(parent context.Context) ([]byte, error) { ctx, cancel := context.WithTimeout(parent, 2*time.Second) defer cancel() return p.client.ReceiveMedia(ctx) } func qualificationSourceVideoPackets(t *testing.T, key []byte, frame uint32, encoded []byte) [][]byte { t.Helper() total := 8 + len(encoded) shardCount := (total + apolloVideoShardPayloadSize - 1) / apolloVideoShardPayloadSize if shardCount > apolloVideoMaximumBlocks*apolloVideoMaximumDataShards { t.Fatal("qualification encoded frame exceeds source-shaped Apollo bound") } combined := make([]byte, shardCount*apolloVideoShardPayloadSize) combined[0], combined[3] = 0x01, 0x01 lastPayloadLength := total - (shardCount-1)*apolloVideoShardPayloadSize binary.LittleEndian.PutUint16(combined[4:6], uint16(lastPayloadLength)) copy(combined[8:], encoded) lastBlock := (shardCount - 1) / apolloVideoMaximumDataShards packets := make([][]byte, 0, shardCount) for globalIndex := 0; globalIndex < shardCount; { block := globalIndex / apolloVideoMaximumDataShards dataShards := min(apolloVideoMaximumDataShards, shardCount-globalIndex) for shardIndex := 0; shardIndex < dataShards; shardIndex++ { flags := byte(0x01) if shardIndex == 0 { flags |= 0x04 } if shardIndex == dataShards-1 { flags |= 0x02 } offset := globalIndex * apolloVideoShardPayloadSize raw := sourceShapedVideoRaw( frame, uint16(frame*1024+uint32(globalIndex)), frame*1024+uint32(globalIndex), flags, dataShards, 0, shardIndex, combined[offset:offset+apolloVideoShardPayloadSize], ) raw[27] = byte(block<<4 | lastBlock<<6) packets = append(packets, sourceEncryptVideoRaw(t, key, raw, qualificationVideoIV(frame, globalIndex))) globalIndex++ } } return packets } func qualificationVideoIV(frame uint32, shard int) string { return fmt.Sprintf("%06x%04xQV", frame&0xffffff, shard&0xffff) } func qualificationProductionPathSmoke(t *testing.T, profile qualificationMediaProfile) qualificationPathTrace { t.Helper() path := newQualificationPath(t, profile, profile.BitrateKbps) defer path.Close() trace, _, err := path.traverse(t, qualificationPayload(profile)) if err != nil { t.Fatal(err) } return trace } func TestQualificationTraversesNativePublicGatewayPath(t *testing.T) { profile := qualificationMediaProfiles()[0] profile.BitrateKbps = 100000 trace := qualificationProductionPathSmoke(t, profile) if !trace.NativeSetup || !trace.NativeOpen || !trace.NativeUDPIngress || !trace.ApolloRecovered || !trace.ProductionQueue || !trace.ProductionMediaLoop || !trace.ProductionPacer || !trace.VerseQUIC || !trace.PublicClientDecode || !trace.PayloadPreserved { t.Fatalf("qualification path skipped production stages: %#v", trace) } } func summarizeQualificationSamples(samples []time.Duration) (qualificationProcessingSummary, error) { if len(samples) == 0 { return qualificationProcessingSummary{}, errors.New("qualification has no processing samples") } var mean, m2 float64 histogram := newQualificationHistogram() for index, sample := range samples { value := float64(sample) delta := value - mean mean += delta / float64(index+1) m2 += delta * (value - mean) observeQualificationHistogram(histogram, sample) } sort.Slice(samples, func(first, second int) bool { return samples[first] < samples[second] }) return qualificationProcessingSummary{ Count: int64(len(samples)), Min: samples[0], Median: qualificationPercentile(samples, 0.50), P90: qualificationPercentile(samples, 0.90), P95: qualificationPercentile(samples, 0.95), P99: qualificationPercentile(samples, 0.99), Max: samples[len(samples)-1], Mean: time.Duration(mean), StandardDeviation: time.Duration(math.Sqrt(m2 / float64(len(samples)))), Histogram: histogram, }, nil } func qualificationPercentile(samples []time.Duration, percentile float64) time.Duration { index := int(math.Ceil(percentile*float64(len(samples)))) - 1 if index < 0 { index = 0 } return samples[index] } func enforceQualificationProcessingGate(summary qualificationProcessingSummary) error { if summary.Count < 1 || summary.P95 > qualificationProcessingLimit { return fmt.Errorf("processing p95 %s exceeds %s", summary.P95, qualificationProcessingLimit) } return nil } func newQualificationHistogram() map[string]int { return map[string]int{ "le_1us": 0, "le_5us": 0, "le_10us": 0, "le_25us": 0, "le_50us": 0, "le_100us": 0, "le_250us": 0, "le_500us": 0, "le_1ms": 0, "le_2ms": 0, "le_5ms": 0, "gt_5ms": 0, } } func observeQualificationHistogram(histogram map[string]int, sample time.Duration) { buckets := []struct { name string limit time.Duration }{ {"le_1us", time.Microsecond}, {"le_5us", 5 * time.Microsecond}, {"le_10us", 10 * time.Microsecond}, {"le_25us", 25 * time.Microsecond}, {"le_50us", 50 * time.Microsecond}, {"le_100us", 100 * time.Microsecond}, {"le_250us", 250 * time.Microsecond}, {"le_500us", 500 * time.Microsecond}, {"le_1ms", time.Millisecond}, {"le_2ms", 2 * time.Millisecond}, {"le_5ms", 5 * time.Millisecond}, } observed := false for _, bucket := range buckets { if sample <= bucket.limit { histogram[bucket.name]++ observed = true } } if !observed { histogram["gt_5ms"]++ } } func runQualificationImpairment(t *testing.T, profile qualificationImpairmentProfile, media qualificationMediaProfile, packetCount int, rawPath string) (qualificationImpairmentObservation, error) { t.Helper() if packetCount < 1 || packetCount > qualificationImpairmentMaxPackets || media.PacketBytes < 1 || media.PacketBytes > 1179 || !qualificationKnownImpairment(profile) { return qualificationImpairmentObservation{}, errors.New("qualification impairment bounds invalid") } file, err := os.OpenFile(rawPath, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o640) if err != nil { return qualificationImpairmentObservation{}, err } compressed := gzip.NewWriter(file) buffered := bufio.NewWriter(compressed) closed := false defer func() { if !closed { _ = buffered.Flush() _ = compressed.Close() _ = file.Close() } }() if _, err := buffered.WriteString("source_sequence,sent_ns,delivered_ns,processing_ns,outcome,delivery_order,bytes,queue_packets\n"); err != nil { return qualificationImpairmentObservation{}, err } type scheduledPacket struct { index int target time.Duration } type rawSample struct { sent, delivered time.Duration processing time.Duration outcome string deliveryOrder int bytes int queuePackets int } type receivedPacket struct { index, deliveryOrder, queuePackets int deliveredAt time.Time processing time.Duration } state := qualificationImpairmentSeed random := func() uint64 { state ^= state << 13 state ^= state >> 7 state ^= state << 17 return state } spacing := time.Duration(int64(time.Second) * int64(media.PacketBytes) * 8 / (media.BitrateKbps * 1000)) if spacing < time.Nanosecond { spacing = time.Nanosecond } observation := qualificationImpairmentObservation{ Profile: profile.Name, MediaProfile: media.Name, Seed: qualificationImpairmentSeed, Sent: packetCount, ConfiguredRTT: profile.RTT, ConfiguredJitter: profile.Jitter, ConfiguredLossPercent: profile.LossPercent, ConfiguredReorder: profile.Reorder, ConfiguredCapacitySteps: append([]int(nil), profile.CapacitySteps...), RTTSource: "apollo_enet_acknowledge", } path := newQualificationImpairedPath(t, media, qualificationMediaPacerKbps(media, 0), profile) defer path.Close() payload := qualificationPayload(media) rawSamples := make([]rawSample, packetCount) jobs := make([]scheduledPacket, 0, packetCount) jobByIndex := make([]int, packetCount) for index := range jobByIndex { jobByIndex[index] = -1 } for index := 0; index < packetCount; index++ { jitter := time.Duration(0) if profile.Jitter > 0 { width := uint64(profile.Jitter*2 + 1) jitter = time.Duration(random()%width) - profile.Jitter } rawSamples[index].sent = time.Duration(index) * spacing if float64(random()%10_000) < profile.LossPercent*100 { rawSamples[index].outcome = "injected_dropped" observation.InjectedDropped++ continue } delay := max(profile.RTT/2+jitter, 0) jobByIndex[index] = len(jobs) jobs = append(jobs, scheduledPacket{index: index, target: time.Duration(index)*spacing + delay}) rawSamples[index].outcome = "traversal_dropped" } var jitterMean, jitterM2 float64 var jitterSamples int for _, packet := range jobs { applied := float64(packet.target - time.Duration(packet.index)*spacing - profile.RTT/2) jitterSamples++ delta := applied - jitterMean jitterMean += delta / float64(jitterSamples) jitterM2 += delta * (applied - jitterMean) } if jitterSamples > 0 { observation.AppliedJitter = time.Duration(math.Sqrt(jitterM2 / float64(jitterSamples))) } if profile.Reorder { for index := 18; index+1 < packetCount; index += 20 { first, second := jobByIndex[index], jobByIndex[index+1] if first < 0 || second < 0 { continue } jobs[first], jobs[second] = jobs[second], jobs[first] observation.InjectedReordered++ } } started := time.Now() beforeMetrics := path.server.Metrics() beforeIngress := path.session.mediaIngress.Load() beforeRecovered := path.session.mediaRecovered.Load() beforeEnqueued := path.session.mediaEnqueued.Load() beforeProviderDrops := path.session.mediaDrops.Load() beforeSourceUDP := path.sourceUDP.Load() beforePacer := path.server.pacer.reservations.Load() grace := max(2*profile.RTT+2*profile.Jitter, 2*time.Second) lastTarget := time.Duration(packetCount) * spacing if len(jobs) > 0 { lastTarget = jobs[len(jobs)-1].target } receiveCtx, receiveCancel := context.WithDeadline(context.Background(), started.Add(lastTarget+grace)) defer receiveCancel() receivedDone := make(chan struct { packets []receivedPacket err error }, 1) go func() { result := struct { packets []receivedPacket err error }{packets: make([]receivedPacket, 0, len(jobs))} metrics := beforeMetrics seen := make([]bool, packetCount) for len(result.packets) < len(jobs) { recovered, receiveErr := path.receivePayload(receiveCtx) if receiveErr != nil { if errors.Is(receiveErr, context.DeadlineExceeded) || errors.Is(receiveErr, context.Canceled) { break } var timeout net.Error if errors.As(receiveErr, &timeout) && timeout.Timeout() { break } result.err = receiveErr break } if len(recovered) != len(payload) { result.err = errors.New("impaired payload length changed") break } index := int(binary.BigEndian.Uint32(recovered[len(recovered)-4:])) if index < 0 || index >= packetCount || jobByIndex[index] < 0 || seen[index] { result.err = errors.New("impaired payload sequence invalid") break } expected := append([]byte(nil), payload...) binary.BigEndian.PutUint32(expected[len(expected)-4:], uint32(index)) if !bytes.Equal(recovered, expected) { result.err = errors.New("impaired payload bytes changed") break } seen[index] = true current := path.server.Metrics() processing := time.Duration(0) if samples := current.ProcessingSamples - metrics.ProcessingSamples; samples > 0 { processing = time.Duration((current.ProcessingDelayNanos - metrics.ProcessingDelayNanos) / samples) } metrics = current result.packets = append(result.packets, receivedPacket{ index: index, deliveryOrder: len(result.packets) + 1, deliveredAt: time.Now(), processing: processing, queuePackets: len(path.session.video), }) } receivedDone <- result }() stepAt := make(map[int]time.Time, len(profile.CapacitySteps)) emitted := 0 var nextRelease time.Time for _, packet := range jobs { release := qualificationBoundedRelease(started.Add(packet.target), nextRelease, time.Now(), spacing) qualificationWaitUntil(release) nextRelease = release.Add(spacing) if len(profile.CapacitySteps) == 2 { switch { case packet.index >= packetCount*2/3 && stepAt[profile.CapacitySteps[1]].IsZero(): path.server.pacer.setKbps(qualificationMediaPacerKbps(media, profile.CapacitySteps[1])) stepAt[profile.CapacitySteps[1]] = time.Now() case packet.index >= packetCount/3 && stepAt[profile.CapacitySteps[0]].IsZero(): path.server.pacer.setKbps(qualificationMediaPacerKbps(media, profile.CapacitySteps[0])) stepAt[profile.CapacitySteps[0]] = time.Now() } } current := append([]byte(nil), payload...) binary.BigEndian.PutUint32(current[len(current)-4:], uint32(packet.index)) if _, err := path.emit(t, current); err != nil { receiveCancel() return qualificationImpairmentObservation{}, err } emitted++ } received := <-receivedDone receiveCancel() if received.err != nil { return qualificationImpairmentObservation{}, received.err } var deliveries []qualificationDeliverySample var totalLatency, totalJitter, previousLatency time.Duration previousDelivered := -1 for _, packet := range received.packets { sample := &rawSamples[packet.index] sample.delivered = packet.deliveredAt.Sub(started) sample.processing = packet.processing sample.outcome = "delivered" sample.deliveryOrder = packet.deliveryOrder sample.bytes = media.PacketBytes sample.queuePackets = packet.queuePackets latency := sample.delivered - sample.sent totalLatency += latency if previousLatency != 0 { delta := latency - previousLatency if delta < 0 { delta = -delta } totalJitter += delta } previousLatency = latency if previousDelivered >= 0 && packet.index < previousDelivered { 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) observation.Dropped = packetCount - observation.Delivered completedAt := started if observation.Delivered > 0 { completedAt = received.packets[len(received.packets)-1].deliveredAt } afterMetrics := path.server.Metrics() observation.SourceEmitted = emitted observation.SourceDatagrams = path.sourceUDP.Load() - beforeSourceUDP observation.ProviderIngressDatagrams = path.session.mediaIngress.Load() - beforeIngress observation.ProviderRecovered = int(path.session.mediaRecovered.Load() - beforeRecovered) observation.ProviderEnqueued = int(path.session.mediaEnqueued.Load() - beforeEnqueued) observation.ProviderQueueReplaced = int(path.session.mediaDrops.Load() - beforeProviderDrops) observation.GatewayForwarded = int(afterMetrics.ProcessingSamples - beforeMetrics.ProcessingSamples) observation.QUICSent = observation.GatewayForwarded observation.QUICDatagramsSent = int(afterMetrics.MediaPackets - beforeMetrics.MediaPackets) observation.ProviderFECDropped = max(observation.SourceEmitted-observation.ProviderRecovered, 0) observation.ProviderEnqueueDropped = max(observation.ProviderRecovered-observation.ProviderEnqueued, 0) observation.GatewayDropped = max(observation.ProviderEnqueued-observation.ProviderQueueReplaced-observation.GatewayForwarded, 0) observation.QUICSendDropped = max(observation.GatewayForwarded-observation.QUICSent, 0) observation.ClientDeliveryDropped = max(observation.QUICSent-observation.Delivered, 0) observation.UnexplainedDropped = observation.Dropped - observation.InjectedDropped - observation.ProviderFECDropped - observation.ProviderEnqueueDropped - observation.ProviderQueueReplaced - observation.GatewayDropped - observation.QUICSendDropped - observation.ClientDeliveryDropped if observation.Delivered > 0 && (path.session.mediaIngress.Load() <= beforeIngress || path.session.mediaRecovered.Load() <= beforeRecovered || path.session.mediaEnqueued.Load() <= beforeEnqueued || afterMetrics.ProcessingSamples <= beforeMetrics.ProcessingSamples || path.server.pacer.reservations.Load() <= beforePacer || afterMetrics.MediaPackets <= beforeMetrics.MediaPackets) { return qualificationImpairmentObservation{}, errors.New("impaired traffic bypassed a production gateway stage") } observation.ObservedRTT = path.session.Telemetry().ControlRTT for index, sample := range rawSamples { if _, err := fmt.Fprintf(buffered, "%d,%d,%d,%d,%s,%d,%d,%d\n", index, sample.sent.Nanoseconds(), sample.delivered.Nanoseconds(), sample.processing.Nanoseconds(), sample.outcome, sample.deliveryOrder, sample.bytes, sample.queuePackets); err != nil { return qualificationImpairmentObservation{}, err } } if err := buffered.Flush(); err != nil { return qualificationImpairmentObservation{}, err } if err := compressed.Close(); err != nil { return qualificationImpairmentObservation{}, err } if err := file.Close(); err != nil { return qualificationImpairmentObservation{}, err } closed = true if observation.Delivered > 0 { observation.ObservedLatency = totalLatency / time.Duration(observation.Delivered) if observation.Delivered > 1 { observation.ObservedJitter = totalJitter / time.Duration(observation.Delivered-1) } observation.ObservedThroughputKbps = float64(observation.Delivered*media.PacketBytes*8) / completedAt.Sub(started).Seconds() / 1000 } observation.ObservedLossPercent = float64(observation.Dropped) * 100 / float64(packetCount) observation.ObservedReorderPercent = float64(observation.ObservedOutOfOrder) * 100 / float64(packetCount) if packetCount >= qualificationImpairmentPacketCount { capacityFactor := 1.0 if len(profile.CapacitySteps) > 0 { rates := 1.0 for _, reduction := range profile.CapacitySteps { rates += float64(100-reduction) / 100 } capacityFactor = rates / float64(len(profile.CapacitySteps)+1) } expectedThroughput := float64(media.BitrateKbps) * (1 - profile.LossPercent/100) * capacityFactor lowerThroughput := expectedThroughput * 0.90 upperThroughput := expectedThroughput * 1.05 if observation.ObservedThroughputKbps < lowerThroughput || observation.ObservedThroughputKbps > upperThroughput { return qualificationImpairmentObservation{}, fmt.Errorf( "observed throughput %.2f outside [%.2f,%.2f] over %s with %d delivered/%d injected drops and steps %v", observation.ObservedThroughputKbps, lowerThroughput, upperThroughput, completedAt.Sub(started), observation.Delivered, observation.InjectedDropped, stepAt, ) } } observation.RawSamples = filepath.Base(rawPath) observation.RawSamplesSHA256, observation.RawSamplesBytes, err = qualificationFileSHA256(rawPath) if err != nil { 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, }) 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) } } if observation.Delivered+observation.Dropped != observation.Sent || observation.UnexplainedDropped != 0 || observation.MaxQueuePackets > qualificationImpairmentQueuePackets { return qualificationImpairmentObservation{}, errors.New("qualification impairment accounting invalid") } return observation, nil } func qualificationKnownImpairment(profile qualificationImpairmentProfile) bool { for _, known := range qualificationImpairmentProfiles() { if profile.Name != known.Name || profile.RTT != known.RTT || profile.Jitter != known.Jitter || profile.LossPercent != known.LossPercent || profile.Reorder != known.Reorder || len(profile.CapacitySteps) != len(known.CapacitySteps) { continue } match := true for index := range known.CapacitySteps { match = match && profile.CapacitySteps[index] == known.CapacitySteps[index] } if match { return true } } return false } func qualificationDeliveriesAfter(deliveries []qualificationDeliverySample, start time.Time) []qualificationDeliverySample { index := sort.Search(len(deliveries), func(index int) bool { return !deliveries[index].At.Before(start) }) return deliveries[index:] } func qualificationMeasuredConvergence(deliveries []qualificationDeliverySample, start time.Time, targetBytesPerSecond int64) time.Duration { const window = 250 * time.Millisecond const requiredWindows = 4 consecutive := 0 for offset := time.Duration(0); offset <= 10*time.Second; offset += window { windowStart := start.Add(offset) var total int64 for _, delivery := range deliveries { if !delivery.At.Before(windowStart) && delivery.At.Before(windowStart.Add(window)) { total += delivery.Bytes } } rate := total * int64(time.Second) / int64(window) if rate >= targetBytesPerSecond*90/100 && rate <= targetBytesPerSecond*105/100 { consecutive++ if consecutive == requiredWindows { return offset + window } } else { consecutive = 0 } if len(deliveries) > 0 && windowStart.After(deliveries[len(deliveries)-1].At) { break } } return 11 * time.Second } func qualificationMaximumDeliveryBytes(deliveries []qualificationDeliverySample, window time.Duration) int64 { var maximum, total int64 for first, last := 0, 0; first < len(deliveries); first++ { for last < len(deliveries) && deliveries[last].At.Sub(deliveries[first].At) <= window { total += deliveries[last].Bytes last++ } if total > maximum { maximum = total } total -= deliveries[first].Bytes } return maximum } func runQualificationProcessing(t *testing.T, profile qualificationMediaProfile, rawPath string) (qualificationProcessingSummary, error) { t.Helper() payload := qualificationFramePayload(profile, 0) if len(payload) < 4 { return qualificationProcessingSummary{}, errors.New("qualification payload too small") } path := newQualificationProcessingPath(t, profile, qualificationFramePacerKbps(profile)) defer path.Close() if err := runQualificationProcessWarmup(t, path, profile, payload); err != nil { return qualificationProcessingSummary{}, err } resourcePath := strings.TrimSuffix(rawPath, ".csv.gz") + "-resources.csv.gz" if err := path.process.startRecording(rawPath, resourcePath); err != nil { return qualificationProcessingSummary{}, err } before, err := path.process.snapshot() if err != nil { return qualificationProcessingSummary{}, err } fps := qualificationFrameRate(profile) targetFrames := profile.Duration.Nanoseconds() * int64(fps) / int64(time.Second) if targetFrames < 1 { return qualificationProcessingSummary{}, errors.New("qualification duration produces no complete frames") } spacing := time.Second / time.Duration(fps) var targetBytes int64 var expectedDatagrams uint64 minimumFrameBytes, maximumFrameBytes := maxCompleteFrameBytes, 0 for index := int64(0); index < targetFrames; index++ { size := qualificationFrameSize(profile, index) targetBytes += int64(size) expectedDatagrams += uint64((size + frameV2PayloadSize - 1) / frameV2PayloadSize) minimumFrameBytes = min(minimumFrameBytes, size) maximumFrameBytes = max(maximumFrameBytes, size) } maximumDuration := profile.Duration*105/100 + 250*time.Millisecond started := time.Now() receiveCtx, receiveCancel := context.WithDeadline(context.Background(), started.Add(maximumDuration+2*time.Second)) defer receiveCancel() receivedDone := make(chan error, 1) go func() { for index := int64(0); index < targetFrames; index++ { recovered, receiveErr := path.receivePayload(receiveCtx) if receiveErr != nil { receivedDone <- receiveErr return } expected := qualificationFramePayload(profile, index) if len(recovered) != len(expected) { receivedDone <- fmt.Errorf( "qualification processing payload length = %d, want %d at sequence %d", len(recovered), len(expected), index, ) return } if !bytes.Equal(recovered, expected) { receivedDone <- fmt.Errorf("qualification processing payload bytes changed at sequence %d", index) return } } receivedDone <- nil }() var processed int64 var nextRelease time.Time payloadDigest := sha256.New() for processed < targetFrames { select { case receiveErr := <-receivedDone: if receiveErr != nil { snapshot, _ := path.process.snapshot() return qualificationProcessingSummary{}, fmt.Errorf("%w; gateway snapshot=%#v", receiveErr, snapshot) } return qualificationProcessingSummary{}, errors.New("qualification processing receiver ended early") default: } release := qualificationBoundedRelease( started.Add(time.Duration(processed)*spacing), nextRelease, time.Now(), spacing, ) qualificationWaitUntil(release) nextRelease = release.Add(spacing) current := qualificationFramePayload(profile, processed) _, _ = payloadDigest.Write(current) if _, err := path.emit(t, current); err != nil { return qualificationProcessingSummary{}, err } processed++ } if err := <-receivedDone; err != nil { snapshot, _ := path.process.snapshot() return qualificationProcessingSummary{}, fmt.Errorf("%w; gateway snapshot=%#v", err, snapshot) } qualificationWaitUntil(started.Add(profile.Duration)) actualDuration := time.Since(started) record, err := path.process.stopRecording() if err != nil { return qualificationProcessingSummary{}, err } after, err := path.process.snapshot() if err != nil { return qualificationProcessingSummary{}, err } if after.NativeSetups != 1 || after.NativeOpens != 1 || after.MediaRecovered-before.MediaRecovered != uint64(processed) || after.MediaEnqueued-before.MediaEnqueued != uint64(processed) || after.MediaDrops != before.MediaDrops || after.MediaQueueMaximum > nativeApolloVideoQueuePackets || after.MediaQueueMaximumBytes > nativeApolloVideoQueueBytes || after.Metrics.ProcessingSamples-before.Metrics.ProcessingSamples != uint64(processed) || after.PacerReservations <= before.PacerReservations || after.Metrics.MediaPackets-before.Metrics.MediaPackets != expectedDatagrams { return qualificationProcessingSummary{}, fmt.Errorf("qualification subprocess bypassed a production stage: before=%#v after=%#v processed=%d", before, after, processed) } samples, err := readQualificationProcessingSamples(rawPath, record.Count) if err != nil { return qualificationProcessingSummary{}, err } summary, err := summarizeQualificationSamples(samples) if err != nil { return qualificationProcessingSummary{}, err } sum, size, err := qualificationFileSHA256(rawPath) if err != nil { return qualificationProcessingSummary{}, err } summary.Profile = profile.Name summary.Codec = profile.Codec summary.ConfiguredBitrateKbps = profile.BitrateKbps summary.ObservedBitrateKbps = float64(targetBytes*8) / actualDuration.Seconds() / 1000 summary.ConfiguredFPS = fps summary.ObservedFPS = float64(processed) / actualDuration.Seconds() summary.PayloadBytes = targetBytes summary.MinimumFrameBytes = minimumFrameBytes summary.MaximumFrameBytes = maximumFrameBytes summary.Warmup = profile.Warmup summary.ConfiguredDuration = profile.Duration summary.ActualDuration = actualDuration summary.ClockOverhead = record.ClockOverhead summary.ClockMethod = record.ClockMethod summary.PayloadSHA256 = fmt.Sprintf("%x", payloadDigest.Sum(nil)) summary.RawSamples = filepath.Base(rawPath) summary.RawSamplesSHA256 = sum summary.RawSamplesBytes = size resourceSum, resourceSize, err := qualificationFileSHA256(resourcePath) if err != nil { return qualificationProcessingSummary{}, err } summary.RawResources = filepath.Base(resourcePath) summary.RawResourcesSHA256 = resourceSum summary.RawResourcesBytes = resourceSize summary.ResourceSamples = record.ResourceSamples summary.CPUScope = qualificationGatewayCPUScope summary.ResourceMethod = qualificationResourceMethod summary.CPUSeconds = record.CPUSeconds summary.PeakHeapBytes = record.PeakHeapBytes summary.PeakGoroutines = record.PeakGoroutines summary.AllocatedObjects = record.AllocatedObjects summary.AllocatedBytes = record.AllocatedBytes if summary.CPUSeconds < 0 || summary.ClockOverhead <= 0 || summary.ClockMethod == "" { return qualificationProcessingSummary{}, errors.New("process CPU usage unavailable") } if actualDuration < profile.Duration || actualDuration > maximumDuration { return qualificationProcessingSummary{}, fmt.Errorf("wall-clock duration %s outside [%s,%s]", actualDuration, profile.Duration, maximumDuration) } if summary.ObservedBitrateKbps < float64(profile.BitrateKbps)*0.95 || summary.ObservedBitrateKbps > float64(profile.BitrateKbps)*1.05 { return qualificationProcessingSummary{}, fmt.Errorf("observed bitrate %.2f outside profile bounds for %d", summary.ObservedBitrateKbps, profile.BitrateKbps) } if err := enforceQualificationProcessingGate(summary); err != nil { return qualificationProcessingSummary{}, err } return summary, nil } func readQualificationProcessingSamples(path string, expected int) ([]time.Duration, error) { file, err := os.Open(path) if err != nil { return nil, err } defer file.Close() compressed, err := gzip.NewReader(file) if err != nil { return nil, err } defer compressed.Close() scanner := bufio.NewScanner(compressed) if !scanner.Scan() || scanner.Text() != "elapsed_ns,queue_ns,processing_ns,pacing_ns" { return nil, errors.New("qualification processing header invalid") } samples := make([]time.Duration, 0, expected) for scanner.Scan() { fields := strings.Split(scanner.Text(), ",") if len(fields) != 4 { return nil, errors.New("qualification processing row invalid") } value, err := strconv.ParseInt(fields[2], 10, 64) if err != nil || value < 0 { return nil, errors.New("qualification processing sample invalid") } samples = append(samples, time.Duration(value)) } if err := scanner.Err(); err != nil { return nil, err } if len(samples) != expected { return nil, fmt.Errorf("qualification processing rows=%d want=%d", len(samples), expected) } return samples, nil } func runQualificationProcessWarmup(t *testing.T, path *qualificationPath, profile qualificationMediaProfile, payload []byte) error { t.Helper() started := time.Now() for time.Since(started) < profile.Warmup { if _, err := path.emit(t, payload); err != nil { return err } recovered, err := path.receivePayload(context.Background()) if err != nil || !bytes.Equal(recovered, payload) { return errors.New("qualification warmup payload integrity failure") } } return nil } func runQualificationWarmup(t *testing.T, path *qualificationPath, profile qualificationMediaProfile, payload []byte) error { t.Helper() started := time.Now() for time.Since(started) < profile.Warmup { trace, _, err := path.traverse(t, payload) if err != nil { return err } if !trace.PayloadPreserved { return errors.New("qualification warmup payload integrity failure") } } return nil } func qualificationRuntimeSample(started time.Time) qualificationResourceSample { var usage syscall.Rusage cpuSeconds := -1.0 if syscall.Getrusage(syscall.RUSAGE_SELF, &usage) == nil { cpuSeconds = float64(usage.Utime.Sec+usage.Stime.Sec) + float64(usage.Utime.Usec+usage.Stime.Usec)/1_000_000 } samples := []runtimemetrics.Sample{ {Name: "/memory/classes/heap/objects:bytes"}, {Name: "/sched/goroutines:goroutines"}, {Name: "/gc/heap/allocs:objects"}, {Name: "/gc/heap/allocs:bytes"}, } runtimemetrics.Read(samples) return qualificationResourceSample{ Elapsed: time.Since(started), CPUSeconds: cpuSeconds, HeapBytes: samples[0].Value.Uint64(), Goroutines: int(samples[1].Value.Uint64()), AllocatedObjects: samples[2].Value.Uint64(), AllocatedBytes: samples[3].Value.Uint64(), } } func writeQualificationResourceSamples(path string, samples []qualificationResourceSample) error { file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o640) if err != nil { return err } compressed := gzip.NewWriter(file) buffered := bufio.NewWriter(compressed) if _, err = buffered.WriteString("elapsed_ns,cpu_seconds,heap_object_bytes,goroutines,allocated_objects,allocated_bytes\n"); err == nil { for _, sample := range samples { if _, err = fmt.Fprintf(buffered, "%d,%.9f,%d,%d,%d,%d\n", sample.Elapsed.Nanoseconds(), sample.CPUSeconds, sample.HeapBytes, sample.Goroutines, sample.AllocatedObjects, sample.AllocatedBytes); err != nil { break } } } if flushErr := buffered.Flush(); err == nil { err = flushErr } if closeErr := compressed.Close(); err == nil { err = closeErr } if closeErr := file.Close(); err == nil { err = closeErr } return err } func qualificationClockOverhead() time.Duration { const readsPerBatch = 100 samples := make([]time.Duration, 1000) var observed time.Time for index := range samples { started := time.Now() for range readsPerBatch { observed = time.Now() } samples[index] = time.Since(started) / readsPerBatch } runtime.KeepAlive(observed) summary, _ := summarizeQualificationSamples(samples) return summary.Median } func qualificationFileSHA256(path string) (string, int64, error) { file, err := os.Open(path) if err != nil { return "", 0, err } defer file.Close() hash := sha256.New() size, err := io.Copy(hash, file) if err != nil { return "", 0, err } return hex.EncodeToString(hash.Sum(nil)), size, nil } type qualificationFlowDelivery struct { at time.Time flow string bytes int64 } func runQualificationFleetStage(t *testing.T, fleet *qualificationFleet, profile qualificationMediaProfile, duration time.Duration) ([]qualificationFlowDelivery, error) { t.Helper() if duration <= 0 { return nil, errors.New("qualification fleet duration must be positive") } end := time.Now().Add(duration) payload := qualificationPayload(profile) ctx, cancel := context.WithCancel(context.Background()) defer cancel() deliveries := make([]qualificationFlowDelivery, 0, int(duration/time.Millisecond)) var wait sync.WaitGroup var mu sync.Mutex var firstErr error for flowIndex, path := range fleet.paths { wait.Add(1) go func(flowIndex int, path *qualificationPath) { defer wait.Done() var sequence uint32 for time.Now().Before(end) && ctx.Err() == nil { current := append([]byte(nil), payload...) binary.BigEndian.PutUint32(current[len(current)-4:], uint32(flowIndex)<<24|sequence) sequence++ trace, _, err := path.traverse(t, current) if err == nil && (!trace.NativeUDPIngress || !trace.ApolloRecovered || !trace.ProductionQueue || !trace.ProductionMediaLoop || !trace.ProductionPacer || !trace.VerseQUIC || !trace.PublicClientDecode || !trace.PayloadPreserved) { err = errors.New("fairness traffic bypassed production gateway traversal") } if err != nil { mu.Lock() if firstErr == nil { firstErr = err cancel() } mu.Unlock() return } mu.Lock() deliveries = append(deliveries, qualificationFlowDelivery{ at: time.Now(), flow: path.flow, bytes: int64(len(current) + frameHeaderSize), }) mu.Unlock() } }(flowIndex, path) } wait.Wait() if firstErr != nil { return nil, firstErr } sort.Slice(deliveries, func(first, second int) bool { return deliveries[first].at.Before(deliveries[second].at) }) return deliveries, nil } func qualificationPacerEvidence(t *testing.T, rawPath string, baselineDuration, stepDuration time.Duration) (qualificationFairnessEvidence, error) { t.Helper() flows := []string{"one", "two", "three", "four", "five", "six", "seven", "eight"} profile := qualificationMediaProfile{Name: "fairness-h264", Codec: "h264", BitrateKbps: 8000, PacketBytes: 1000} fleet := newQualificationFleet(t, len(flows), profile, 8000) defer fleet.Close() for index, path := range fleet.paths { path.flow = flows[index] if _, _, err := path.traverse(t, qualificationPayload(profile)); err != nil { return qualificationFairnessEvidence{}, err } } start := time.Now() baseline, err := runQualificationFleetStage(t, fleet, profile, baselineDuration) if err != nil { return qualificationFairnessEvidence{}, err } allDeliveries := append([]qualificationFlowDelivery(nil), baseline...) evidence := qualificationFairnessEvidence{ Evaluation: baselineDuration, PerFlowBytes: make(map[string]int64, len(flows)), ShareError: make(map[string]float64, len(flows)), } var total, squares float64 for _, delivery := range baseline { evidence.PerFlowBytes[delivery.flow] += delivery.bytes } for _, flow := range flows { total += float64(evidence.PerFlowBytes[flow]) squares += float64(evidence.PerFlowBytes[flow]) * float64(evidence.PerFlowBytes[flow]) } target := total / float64(len(flows)) for _, flow := range flows { evidence.ShareError[flow] = math.Abs(float64(evidence.PerFlowBytes[flow])-target) / target if evidence.ShareError[flow] > 0.10 { return qualificationFairnessEvidence{}, fmt.Errorf("flow %s share error %.4f", flow, evidence.ShareError[flow]) } } evidence.JainIndex = total * total / (float64(len(flows)) * squares) for _, step := range []struct { reduction int kbps int64 cap int64 }{ {25, 6000, 750_000}, {50, 4000, 500_000}, } { fleet.server.pacer.setKbps(step.kbps) stepStart := time.Now() deliveries, runErr := runQualificationFleetStage(t, fleet, profile, stepDuration) if runErr != nil { return qualificationFairnessEvidence{}, runErr } allDeliveries = append(allDeliveries, deliveries...) convergence := qualificationPacerConvergence(deliveries, stepStart, flows, step.cap) maximum := qualificationMaximumFiveSecondBytes(deliveries) if stepDuration >= 10*time.Second && (convergence > 10*time.Second || maximum > step.cap*5*105/100) { return qualificationFairnessEvidence{}, fmt.Errorf("capacity step %d failed convergence=%s five-second=%d", step.reduction, convergence, maximum) } evidence.CapacitySteps = append(evidence.CapacitySteps, qualificationCapacityStep{ ReductionPercent: step.reduction, Convergence: convergence, MaximumFiveSecond: maximum, FiveSecondCap: step.cap * 5, }) } end := start.Add(baselineDuration + 2*stepDuration) evidence.Series = qualificationFairnessSeriesFor(allDeliveries, start, end, flows) if err := writeQualificationPacerSamples(rawPath, allDeliveries, start); err != nil { return qualificationFairnessEvidence{}, err } evidence.RawSamples = filepath.Base(rawPath) evidence.RawSamplesSHA256, evidence.RawSamplesBytes, err = qualificationFileSHA256(rawPath) if err != nil { return qualificationFairnessEvidence{}, err } return evidence, nil } func qualificationPacerConvergence(deliveries []qualificationFlowDelivery, start time.Time, flows []string, targetBytesPerSecond int64) time.Duration { consecutive := 0 for second := time.Duration(0); second < 10*time.Second; second += time.Second { windowStart := start.Add(second) perFlow := make(map[string]int64, len(flows)) var aggregate int64 for _, delivery := range deliveries { if !delivery.at.Before(windowStart) && delivery.at.Before(windowStart.Add(time.Second)) { perFlow[delivery.flow] += delivery.bytes aggregate += delivery.bytes } } if aggregate < targetBytesPerSecond*90/100 || aggregate > targetBytesPerSecond*105/100 { consecutive = 0 continue } targetFlow := targetBytesPerSecond / int64(len(flows)) converged := true for _, flow := range flows { converged = converged && perFlow[flow] >= targetFlow*90/100 && perFlow[flow] <= targetFlow*110/100 } if converged { consecutive++ if consecutive == 2 { return second + time.Second } } else { consecutive = 0 } } return 11 * time.Second } func qualificationMaximumFiveSecondBytes(deliveries []qualificationFlowDelivery) int64 { sort.Slice(deliveries, func(first, second int) bool { return deliveries[first].at.Before(deliveries[second].at) }) var maximum, total int64 for first, last := 0, 0; first < len(deliveries); first++ { for last < len(deliveries) && deliveries[last].at.Sub(deliveries[first].at) <= 5*time.Second { total += deliveries[last].bytes last++ } if total > maximum { maximum = total } total -= deliveries[first].bytes } return maximum } func qualificationFairnessSeriesFor(deliveries []qualificationFlowDelivery, start, end time.Time, flows []string) []qualificationFairnessSeries { series := make([]qualificationFairnessSeries, 0, int(end.Sub(start)/time.Second)) for windowStart := start; windowStart.Before(end); windowStart = windowStart.Add(time.Second) { sample := qualificationFairnessSeries{ Elapsed: windowStart.Sub(start), PerFlowBytes: make(map[string]int64, len(flows)), } for _, delivery := range deliveries { if !delivery.at.Before(windowStart) && delivery.at.Before(windowStart.Add(time.Second)) { sample.PerFlowBytes[delivery.flow] += delivery.bytes sample.AggregateBytes += delivery.bytes } } series = append(series, sample) } return series } func writeQualificationPacerSamples(path string, deliveries []qualificationFlowDelivery, start time.Time) error { file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o640) if err != nil { return err } compressed := gzip.NewWriter(file) buffered := bufio.NewWriter(compressed) if _, err = buffered.WriteString("elapsed_ns,flow,bytes\n"); err == nil { for _, delivery := range deliveries { if _, err = fmt.Fprintf(buffered, "%d,%s,%d\n", delivery.at.Sub(start).Nanoseconds(), delivery.flow, delivery.bytes); err != nil { break } } } if flushErr := buffered.Flush(); err == nil { err = flushErr } if closeErr := compressed.Close(); err == nil { err = closeErr } if closeErr := file.Close(); err == nil { err = closeErr } return err } func qualificationProtocolVersion() (string, error) { version := os.Getenv("VERSEVDI_QUALIFICATION_PROTOCOL_VERSION") valid, err := regexp.MatchString(`^v[0-9]+\.[0-9]+\.[0-9]+-[0-9A-Za-z]+(?:[.-][0-9A-Za-z]+)*$`, version) if err != nil || !valid { return "", errors.New("VERSEVDI_QUALIFICATION_PROTOCOL_VERSION must be an immutable prerelease tag") } return version, nil } func qualificationCandidateCommit() (string, error) { commit := os.Getenv("VERSEVDI_QUALIFICATION_COMMIT") decoded, err := hex.DecodeString(commit) if err != nil || len(decoded) != 20 || strings.ToLower(commit) != commit { return "", errors.New("VERSEVDI_QUALIFICATION_COMMIT must be a lowercase full SHA-1") } return commit, nil } func qualificationToolVersions() (map[string]string, error) { versions := map[string]string{"qualification": qualificationToolVersion, "go": runtime.Version()} info, ok := debug.ReadBuildInfo() if ok { for _, dependency := range info.Deps { if dependency.Path == "github.com/quic-go/quic-go" { versions["quic-go"] = dependency.Version break } } } if versions["quic-go"] == "" { module, err := os.ReadFile(filepath.Join("..", "go.mod")) if err != nil { return nil, errors.New("qualification QUIC implementation version unavailable") } match := regexp.MustCompile(`(?m)^\s*github\.com/quic-go/quic-go\s+(v[^\s]+)`).FindSubmatch(module) if len(match) != 2 { return nil, errors.New("qualification QUIC implementation version unavailable") } versions["quic-go"] = string(match[1]) } return versions, nil } func writeQualificationJSON(path string, value any) error { file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o640) if err != nil { return err } encoder := json.NewEncoder(file) encoder.SetIndent("", " ") if err := encoder.Encode(value); err != nil { _ = file.Close() return err } return file.Close() } func TestSection7Qualification(t *testing.T) { output := os.Getenv("VERSEVDI_QUALIFICATION_DIR") if output == "" { t.Skip("set VERSEVDI_QUALIFICATION_DIR to run the 30-minute frozen-candidate qualification") } if err := validateQualificationOutputDir(output); err != nil { t.Fatal(err) } commit, err := qualificationCandidateCommit() if err != nil { t.Fatal(err) } protocolVersion, err := qualificationProtocolVersion() if err != nil { t.Fatal(err) } toolVersions, err := qualificationToolVersions() if err != nil { t.Fatal(err) } if err := os.Mkdir(output, 0o750); err != nil { t.Fatalf("create new qualification evidence directory: %v", err) } started := time.Now().UTC() media := qualificationMediaProfiles() qualificationTraverseProfiles(t, media) manifest := qualificationManifest{ Status: "running", ToolVersion: qualificationToolVersion, Command: fmt.Sprintf( "VERSEVDI_QUALIFICATION_DIR=%s VERSEVDI_QUALIFICATION_COMMIT=%s VERSEVDI_QUALIFICATION_PROTOCOL_VERSION=%s GOWORK=off go test ./gateway -run '^TestSection7Qualification$' -count=1 -timeout 45m -v", output, commit, protocolVersion, ), CandidateCommit: commit, ProtocolVersion: protocolVersion, StartedAt: started.Format(time.RFC3339Nano), GoVersion: runtime.Version(), ToolVersions: toolVersions, OS: runtime.GOOS, Architecture: runtime.GOARCH, Topology: "parent source-shaped encrypted Apollo fixture -> isolated gateway subprocess for processing/resource evidence -> public Verse client decoder; impairment uses the same native recovery/FEC, bounded queue, production pacer, framing, and QUIC path", Direction: "provider_to_client", QueueDiscipline: "ordered fixed-seed source delay queue with one-serialization-interval catch-up, bounded 256-packet native video queue, production equal-tier fair pacer", Evidence: []string{"deterministic source-shaped Apollo recovery", "isolated gateway-process resources", "local real-time production path", "mTLS/QUIC fixture transport", "attributed path impairment", "production fair pacer"}, Deferred: []string{"live Apollo", "macOS client", "physical firewall and packet route", "real encoder fidelity", "multi-host scale"}, } for _, profile := range media { raw := filepath.Join(output, "processing-"+profile.Name+".csv.gz") summary, runErr := runQualificationProcessing(t, profile, raw) if runErr != nil { t.Fatal(runErr) } manifest.Processing = append(manifest.Processing, summary) t.Logf("%s count=%d p95=%s observed=%.2f kbps", profile.Name, summary.Count, summary.P95, summary.ObservedBitrateKbps) } fairness, err := qualificationPacerEvidence(t, filepath.Join(output, "fairness.csv.gz"), 60*time.Second, 10*time.Second) if err != nil { t.Fatal(err) } manifest.Fairness = fairness for _, impairment := range qualificationImpairmentProfiles() { profiles := media[:1] if impairment.Name == "baseline" { profiles = media } for _, profile := range profiles { raw := filepath.Join(output, "impairment-"+impairment.Name+"-"+profile.Name+".csv.gz") observation, runErr := runQualificationImpairment(t, impairment, profile, qualificationImpairmentPacketCount, raw) if runErr != nil { t.Fatal(runErr) } manifest.Impairments = append(manifest.Impairments, observation) } } manifest.Status = "passed" manifest.CompletedAt = time.Now().UTC().Format(time.RFC3339Nano) manifestPath := filepath.Join(output, "manifest.json") if err := writeQualificationJSON(manifestPath, manifest); err != nil { t.Fatal(err) } manifestBytes, err := os.ReadFile(manifestPath) if err != nil { t.Fatal(err) } for _, forbidden := range []string{"private_key", "client_private", "clipboard text", "apollo.test", "provider endpoint", "live interoperability passed"} { if bytes.Contains(bytes.ToLower(manifestBytes), []byte(forbidden)) { t.Fatalf("qualification manifest contains forbidden boundary text %q", forbidden) } } t.Logf("qualification manifest: %s", manifestPath) } func qualificationTraverseProfiles(t *testing.T, profiles []qualificationMediaProfile) { t.Helper() for _, profile := range profiles { trace := qualificationProductionPathSmoke(t, profile) if !trace.NativeSetup || !trace.NativeOpen || !trace.NativeUDPIngress || !trace.ApolloRecovered || !trace.ProductionQueue || !trace.ProductionMediaLoop || !trace.ProductionPacer || !trace.VerseQUIC || !trace.PublicClientDecode || !trace.PayloadPreserved { t.Fatalf("%s production path: %#v", profile.Name, trace) } } }