package gateway import ( "bufio" "bytes" "compress/gzip" "context" "crypto/aes" "crypto/cipher" "crypto/sha256" "crypto/tls" "encoding/binary" "encoding/csv" "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" "testing" "time" protocol "git.sechmachine.io.vn/sechmachine/VerseVDI-Protocol/gen/go/protocol" ) const ( qualificationToolVersion = "versevdi-gateway-qualification/v10" 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" qualificationApolloVideoRateBitsPerSecond = 1_000_000_000 * 80 / 100 qualificationApolloVideoBatchBytes = 64 * 1024 qualificationApolloVideoBatchPackets = 64 qualificationWireTimebase = "monotonic offsets from constrained run start" ) 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 qualificationVideoBatchObservation struct { ProviderFrame uint32 SourcePacket uint64 PacketWithinFrame int StartedAfter time.Duration } 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 expectedSourceUDP atomic.Uint64 closeOnce sync.Once shutdown func() } type qualificationProcessingDiagnostics struct { ProcessedFrames int64 SourceFrames uint32 ExpectedWritesThroughLastFrame uint64 FixtureSentPackets uint64 SourceUDP uint64 BatchHistory []qualificationVideoBatchObservation Gateway qualificationGatewayProcessSnapshot SnapshotError string } 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"` RawWireSamples string `json:"raw_wire_samples,omitempty"` RawWireSamplesSHA256 string `json:"raw_wire_samples_sha256,omitempty"` RawWireSamplesBytes int64 `json:"raw_wire_samples_bytes,omitempty"` RawWireRows int `json:"raw_wire_rows,omitempty"` RawWireDeliveryRows int `json:"raw_wire_delivery_rows,omitempty"` RawWireTransitionRows int `json:"raw_wire_transition_rows,omitempty"` RawWireTimebase string `json:"raw_wire_timebase,omitempty"` } type qualificationCapacityStep struct { ReductionPercent int `json:"reduction_percent"` TransitionAfter time.Duration `json:"transition_after_ns"` Convergence time.Duration `json:"convergence_ns"` MaximumFiveSecond int64 `json:"maximum_five_second_bytes"` FiveSecondCap int64 `json:"five_second_cap_bytes"` RecomputationSource string `json:"recomputation_source"` } 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 qualificationWireFileEvidence struct { Name string SHA256 string Bytes int64 Rows int DeliveryRows int TransitionRows int } 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 videoPaceMu sync.Mutex videoNext time.Time beforeVideoBatch func(context.Context, int) error beforeVideoFirstWrite func(context.Context, int) error observeVideoBatch func(int, time.Time) videoBatchMu sync.Mutex videoBatchEpoch time.Time videoBatchCount uint64 videoBatches [128]qualificationVideoBatchObservation controlImpairmentMu sync.Mutex controlRTT time.Duration controlJitter time.Duration controlRandom uint64 } func (f *qualificationApolloFixture) videoBatchHistory() []qualificationVideoBatchObservation { f.videoBatchMu.Lock() defer f.videoBatchMu.Unlock() count := min(f.videoBatchCount, uint64(len(f.videoBatches))) result := make([]qualificationVideoBatchObservation, 0, count) for index := f.videoBatchCount - count; index < f.videoBatchCount; index++ { result = append(result, f.videoBatches[index%uint64(len(f.videoBatches))]) } return result } func (f *qualificationApolloFixture) recordVideoBatch(providerFrame uint32, sourcePacket uint64, packetWithinFrame int, started time.Time) { f.videoBatchMu.Lock() if f.videoBatchEpoch.IsZero() { f.videoBatchEpoch = started } f.videoBatches[f.videoBatchCount%uint64(len(f.videoBatches))] = qualificationVideoBatchObservation{ ProviderFrame: providerFrame, SourcePacket: sourcePacket, PacketWithinFrame: packetWithinFrame, StartedAfter: started.Sub(f.videoBatchEpoch), } f.videoBatchCount++ f.videoBatchMu.Unlock() } 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, providerFrame uint32, 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() if len(packets) == 0 { return nil } packetsPerMillisecond, batchSize := qualificationApolloVideoPacing(apolloVideoRawPacketSize) if packetsPerMillisecond == 0 || batchSize == 0 { return ErrProviderMalformed } f.videoPaceMu.Lock() defer f.videoPaceMu.Unlock() framePackets := 0 for batchStart := 0; batchStart < len(packets); batchStart += batchSize { batchEnd := min(batchStart+batchSize, len(packets)) sourceStart := f.sentPackets.Load() if err := qualificationWaitContext(ctx, f.videoNext); err != nil { return err } if f.beforeVideoBatch != nil { if err := f.beforeVideoBatch(ctx, framePackets); err != nil { return err } } if f.beforeVideoFirstWrite != nil { if err := f.beforeVideoFirstWrite(ctx, framePackets); err != nil { return err } } firstPacket := packets[batchStart] if len(firstPacket) != len(packets[0]) { return ErrProviderMalformed } if _, err := f.video.WriteToUDP(firstPacket, remote); err != nil { return err } batchStarted := time.Now() f.sentPackets.Add(1) f.recordVideoBatch(providerFrame, sourceStart, framePackets, batchStarted) if f.observeVideoBatch != nil { f.observeVideoBatch(framePackets, batchStarted) } for _, packet := range packets[batchStart+1 : batchEnd] { if len(packet) != len(packets[0]) { return ErrProviderMalformed } if _, err := f.video.WriteToUDP(packet, remote); err != nil { return err } f.sentPackets.Add(1) } currentBatch := batchEnd - batchStart framePackets += currentBatch f.videoNext = batchStarted.Add(qualificationApolloVideoOffset(currentBatch, packetsPerMillisecond)) } return nil } func qualificationApolloVideoPacing(packetBytes int) (packetsPerMillisecond, batchSize int) { if packetBytes <= 0 { return 0, 0 } packetsPerMillisecond = qualificationApolloVideoRateBitsPerSecond / 1000 / packetBytes / 8 batchSize = min(qualificationApolloVideoBatchBytes/packetBytes, qualificationApolloVideoBatchPackets) return packetsPerMillisecond, batchSize } func qualificationApolloVideoOffset(packets, packetsPerMillisecond int) time.Duration { return time.Millisecond * time.Duration(packets) / time.Duration(packetsPerMillisecond) } func qualificationWaitContext(ctx context.Context, due time.Time) error { delay := time.Until(due) if delay <= 0 { select { case <-ctx.Done(): return ctx.Err() default: return nil } } timer := time.NewTimer(delay) defer timer.Stop() select { case <-ctx.Done(): return ctx.Err() case <-timer.C: 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 { return newQualificationPathWithNativeBackend(t, profile, pacerKbps, impairment, NewNativeApolloBackend()) } func newQualificationPathWithNativeBackend(t *testing.T, profile qualificationMediaProfile, pacerKbps int64, impairment *qualificationImpairmentProfile, native *NativeApolloBackend) *qualificationPath { return newQualificationPathWithNativeBackendAndObserver(t, profile, pacerKbps, impairment, native, nil) } func newQualificationPathWithNativeBackendAndObserver(t *testing.T, profile qualificationMediaProfile, pacerKbps int64, impairment *qualificationImpairmentProfile, native *NativeApolloBackend, observer func(mediaTimingObservation)) *qualificationPath { t.Helper() serverTLS, clientTLS := testTLS(t) fixture := newQualificationApolloFixture(t, serverTLS, clientTLS, "qualification-session", profile) if impairment != nil { fixture.setControlImpairment(*impairment) } backend := &qualificationTracingBackend{native: native} 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, mediaObserver: observer, }) 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) p.expectedSourceUDP.Add(uint64(len(packets))) ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() if err := p.fixture.sendVideo(ctx, p.frame, packets); err != nil { return qualificationPathTrace{}, err } p.sourceUDP.Add(uint64(len(packets))) return qualificationPathTrace{}, nil } func (p *qualificationPath) processingDiagnostics(processed int64) qualificationProcessingDiagnostics { diagnostics := qualificationProcessingDiagnostics{ ProcessedFrames: processed, SourceFrames: p.frame, ExpectedWritesThroughLastFrame: p.expectedSourceUDP.Load(), FixtureSentPackets: p.fixture.sentPackets.Load(), SourceUDP: p.sourceUDP.Load(), BatchHistory: p.fixture.videoBatchHistory(), } var err error diagnostics.Gateway, err = p.process.snapshot() if err != nil { diagnostics.SnapshotError = err.Error() } return diagnostics } 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, observe) } 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() wireFragmentsPerUnit := (media.PacketBytes + frameV2PayloadSize - 1) / frameV2PayloadSize wireDeliveryLimit := len(jobs) * wireFragmentsPerUnit wireDeliveries := make([]qualificationDeliverySample, 0, wireDeliveryLimit) wireOverflow := false 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.receivePayloadObserved(receiveCtx, func(receivedAt time.Time, bytes int) { if len(wireDeliveries) >= wireDeliveryLimit { wireOverflow = true return } wireDeliveries = append(wireDeliveries, qualificationDeliverySample{At: receivedAt, Bytes: int64(bytes)}) }) 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 } if wireOverflow { return qualificationImpairmentObservation{}, errors.New("qualification public-wire observation bound exceeded") } 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 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 } wirePath := strings.TrimSuffix(rawPath, ".csv.gz") + "-wire.csv.gz" for _, reduction := range profile.CapacitySteps { step := qualificationCapacityStepObservation(wireDeliveries, stepAt[reduction], media, reduction) step.TransitionAfter = stepAt[reduction].Sub(started) step.RecomputationSource = filepath.Base(wirePath) observation.CapacityStepObservations = append(observation.CapacityStepObservations, step) } if len(profile.CapacitySteps) > 0 { wireEvidence, writeErr := writeQualificationWireSamples( wirePath, started, profile.CapacitySteps, stepAt, wireDeliveries, wireDeliveryLimit+len(profile.CapacitySteps), ) if writeErr != nil { return qualificationImpairmentObservation{}, writeErr } observation.RawWireSamples = wireEvidence.Name observation.RawWireSamplesSHA256 = wireEvidence.SHA256 observation.RawWireSamplesBytes = wireEvidence.Bytes observation.RawWireRows = wireEvidence.Rows observation.RawWireDeliveryRows = wireEvidence.DeliveryRows observation.RawWireTransitionRows = wireEvidence.TransitionRows observation.RawWireTimebase = qualificationWireTimebase } for _, step := range observation.CapacityStepObservations { if packetCount >= qualificationImpairmentPacketCount && (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", step.ReductionPercent, step.Convergence, step.MaximumFiveSecond) } } 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 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 || 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 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 := windowOrigin.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 windowStart.Add(window).Sub(start) } } 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) (summary qualificationProcessingSummary, err 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)) var processed int64 defer func() { if err != nil { diagnostics := path.processingDiagnostics(processed) err = fmt.Errorf("%w; diagnostics=%#v", err, diagnostics) } 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 nextRelease time.Time payloadDigest := sha256.New() for processed < targetFrames { select { case receiveErr := <-receivedDone: if receiveErr != nil { return qualificationProcessingSummary{}, receiveErr } 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 { return qualificationProcessingSummary{}, err } 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 { cpuSeconds := qualificationProcessCPUSeconds() 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 } func writeQualificationWireSamples( path string, epoch time.Time, reductions []int, transitions map[int]time.Time, deliveries []qualificationDeliverySample, maximumRows int, ) (qualificationWireFileEvidence, error) { type wireRecord struct { kind string reduction int after time.Duration bytes int64 } if maximumRows < len(reductions)+len(deliveries) { return qualificationWireFileEvidence{}, errors.New("qualification public-wire row bound exceeded") } records := make([]wireRecord, 0, len(reductions)+len(deliveries)) for _, reduction := range reductions { at := transitions[reduction] if at.IsZero() || at.Before(epoch) { return qualificationWireFileEvidence{}, fmt.Errorf("qualification capacity transition %d missing or before epoch", reduction) } records = append(records, wireRecord{kind: "transition", reduction: reduction, after: at.Sub(epoch)}) } for _, delivery := range deliveries { if delivery.At.Before(epoch) || delivery.Bytes <= 0 { return qualificationWireFileEvidence{}, errors.New("qualification public-wire delivery invalid") } records = append(records, wireRecord{kind: "delivery", after: delivery.At.Sub(epoch), bytes: delivery.Bytes}) } sort.SliceStable(records, func(first, second int) bool { if records[first].after != records[second].after { return records[first].after < records[second].after } return records[first].kind == "transition" && records[second].kind != "transition" }) file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o640) if err != nil { return qualificationWireFileEvidence{}, err } compressed := gzip.NewWriter(file) buffered := bufio.NewWriter(compressed) writer := csv.NewWriter(buffered) closeAll := func() error { writer.Flush() if err := writer.Error(); err != nil { _ = compressed.Close() _ = file.Close() return err } if err := buffered.Flush(); err != nil { _ = compressed.Close() _ = file.Close() return err } if err := compressed.Close(); err != nil { _ = file.Close() return err } return file.Close() } if err := writer.Write([]string{"record_type", "reduction_percent", "transition_after_ns", "received_after_ns", "encoded_bytes"}); err != nil { _ = closeAll() return qualificationWireFileEvidence{}, err } evidence := qualificationWireFileEvidence{Name: filepath.Base(path)} for _, record := range records { var row []string switch record.kind { case "transition": row = []string{"transition", strconv.Itoa(record.reduction), strconv.FormatInt(record.after.Nanoseconds(), 10), "", ""} evidence.TransitionRows++ case "delivery": row = []string{"delivery", "", "", strconv.FormatInt(record.after.Nanoseconds(), 10), strconv.FormatInt(record.bytes, 10)} evidence.DeliveryRows++ default: _ = closeAll() return qualificationWireFileEvidence{}, errors.New("qualification public-wire record type invalid") } if err := writer.Write(row); err != nil { _ = closeAll() return qualificationWireFileEvidence{}, err } } if err := closeAll(); err != nil { return qualificationWireFileEvidence{}, err } evidence.Rows = evidence.TransitionRows + evidence.DeliveryRows evidence.SHA256, evidence.Bytes, err = qualificationFileSHA256(path) return evidence, err } 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 := qualificationFairnessFlows() 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, TransitionAfter: stepStart.Sub(start), Convergence: convergence, MaximumFiveSecond: maximum, FiveSecondCap: step.cap * 5, RecomputationSource: filepath.Base(rawPath), }) } 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 qualificationFairnessFlows() []string { return []string{"one", "two", "three", "four", "five", "six", "seven", "eight"} } 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) } } }