diff --git a/gen/go/protocol/protocol.go b/gen/go/protocol/protocol.go index 67bd812..6717a30 100644 --- a/gen/go/protocol/protocol.go +++ b/gen/go/protocol/protocol.go @@ -13,7 +13,7 @@ import ( "time" ) -const SchemaSHA256 = "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" +const SchemaSHA256 = "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b" const ProtocolVersion = "1.0.0" const CurrentWireVersion = "1" const NMinus1WireVersion = "0" @@ -191,13 +191,14 @@ type GatewayDrain struct { } type GatewayHeartbeat struct { - Version string `json:"version"` - GatewayID string `json:"gateway_id"` - Sequence int64 `json:"sequence"` - ObservedAt string `json:"observed_at"` - ActiveConnections int64 `json:"active_connections"` - EgressKbps int64 `json:"egress_kbps"` - State string `json:"state"` + Version string `json:"version"` + GatewayID string `json:"gateway_id"` + Sequence int64 `json:"sequence"` + ObservedAt string `json:"observed_at"` + ActiveConnections int64 `json:"active_connections"` + EgressKbps int64 `json:"egress_kbps"` + State string `json:"state"` + Telemetry GatewayTelemetry `json:"telemetry"` } type GatewayRegistration struct { @@ -216,6 +217,27 @@ type GatewayRegistration struct { Capabilities CapabilityProfile `json:"capabilities"` } +type GatewayTelemetry struct { + AdmittedSessions int64 `json:"admitted_sessions"` + AdmissionRejects int64 `json:"admission_rejects"` + Reconnects int64 `json:"reconnects"` + DrainTransitions int64 `json:"drain_transitions"` + MediaDrops int64 `json:"media_drops"` + MediaPackets int64 `json:"media_packets"` + MediaBytes int64 `json:"media_bytes"` + QueueDelayMicros int64 `json:"queue_delay_micros"` + ProcessingDelayMicros int64 `json:"processing_delay_micros"` + ProcessingSamples int64 `json:"processing_samples"` + PacingDelayMicros int64 `json:"pacing_delay_micros"` + ProviderErrors int64 `json:"provider_errors"` + InputRejected int64 `json:"input_rejected"` + ControlRttMicros int64 `json:"control_rtt_micros"` + ControlJitterMicros int64 `json:"control_jitter_micros"` + ControlLossPpm int64 `json:"control_loss_ppm"` + PendingReliable int64 `json:"pending_reliable"` + ProviderState string `json:"provider_state"` +} + type GrantReference struct { OpaqueValue string `json:"opaque_value"` ExpiresAt string `json:"expires_at"` @@ -265,25 +287,26 @@ type PageInfo struct { } type ProviderSessionWork struct { - Version string `json:"version"` - SessionID string `json:"session_id"` - GatewayID string `json:"gateway_id"` - ReconnectSequence int64 `json:"reconnect_sequence"` - ExpiresAt string `json:"expires_at"` - ProviderProfile string `json:"provider_profile"` - ProviderIdentity string `json:"provider_identity"` - PolicyVersionID string `json:"policy_version_id"` - ApplicationID string `json:"application_id"` - ClientID string `json:"client_id"` - ManagementHost string `json:"management_host"` - ManagementPort int64 `json:"management_port"` - StreamHost string `json:"stream_host"` - StreamPort int64 `json:"stream_port"` - ClientCertificatePem string `json:"client_certificate_pem"` - ClientPrivateKeyPem string `json:"client_private_key_pem"` - ServerCertificatePem string `json:"server_certificate_pem"` - ClipboardPolicy ClipboardPolicy `json:"clipboard_policy"` - ProviderApplicationTerminationAllowed bool `json:"provider_application_termination_allowed"` + Version string `json:"version"` + SessionID string `json:"session_id"` + GatewayID string `json:"gateway_id"` + ReconnectSequence int64 `json:"reconnect_sequence"` + ExpiresAt string `json:"expires_at"` + ProviderProfile string `json:"provider_profile"` + ProviderIdentity string `json:"provider_identity"` + PolicyVersionID string `json:"policy_version_id"` + StreamPolicy ProviderStreamPolicy `json:"stream_policy"` + ApplicationID string `json:"application_id"` + ClientID string `json:"client_id"` + ManagementHost string `json:"management_host"` + ManagementPort int64 `json:"management_port"` + StreamHost string `json:"stream_host"` + StreamPort int64 `json:"stream_port"` + ClientCertificatePem string `json:"client_certificate_pem"` + ClientPrivateKeyPem string `json:"client_private_key_pem"` + ServerCertificatePem string `json:"server_certificate_pem"` + ClipboardPolicy ClipboardPolicy `json:"clipboard_policy"` + ProviderApplicationTerminationAllowed bool `json:"provider_application_termination_allowed"` } type ProviderState struct { @@ -294,6 +317,15 @@ type ProviderState struct { Channels []string `json:"channels"` } +type ProviderStreamPolicy struct { + ResolutionWidth int64 `json:"resolution_width"` + ResolutionHeight int64 `json:"resolution_height"` + Fps int64 `json:"fps"` + Codec string `json:"codec"` + BitrateKbps int64 `json:"bitrate_kbps"` + AudioEnabled bool `json:"audio_enabled"` +} + type ReauthGrant struct { Token string `json:"token"` Purpose string `json:"purpose"` @@ -2370,6 +2402,12 @@ func (v GatewayHeartbeat) Validate() error { if v.State != "" && !(v.State == "ready" || v.State == "draining" || v.State == "offline") { violations = append(violations, FieldViolation{Field: "state", Code: "invalid_value"}) } + if reflect.DeepEqual(v.Telemetry, GatewayTelemetry{}) { + violations = append(violations, FieldViolation{Field: "telemetry", Code: "required"}) + } + if err := v.Telemetry.Validate(); err != nil { + violations = append(violations, FieldViolation{Field: "telemetry", Code: "invalid_object"}) + } if len(violations) > 0 { return ValidationError{Violations: violations} } @@ -2403,6 +2441,9 @@ func DecodeGatewayHeartbeat(data []byte) (GatewayHeartbeat, error) { if raw, ok := fields["state"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { return value, ValidationError{Violations: []FieldViolation{{Field: "state", Code: "required"}}} } + if raw, ok := fields["telemetry"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "telemetry", Code: "required"}}} + } if raw, ok := fields["version"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { return value, ValidationError{Violations: []FieldViolation{{Field: "version", Code: "required"}}} } @@ -2623,6 +2664,210 @@ func EncodeGatewayRegistration(value GatewayRegistration) ([]byte, error) { return json.Marshal(value) } +func (v GatewayTelemetry) Validate() error { + var violations []FieldViolation + if v.AdmittedSessions != 0 && v.AdmittedSessions < 0 { + violations = append(violations, FieldViolation{Field: "admitted_sessions", Code: "minimum"}) + } + if v.AdmittedSessions > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "admitted_sessions", Code: "maximum"}) + } + if v.AdmissionRejects != 0 && v.AdmissionRejects < 0 { + violations = append(violations, FieldViolation{Field: "admission_rejects", Code: "minimum"}) + } + if v.AdmissionRejects > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "admission_rejects", Code: "maximum"}) + } + if v.Reconnects != 0 && v.Reconnects < 0 { + violations = append(violations, FieldViolation{Field: "reconnects", Code: "minimum"}) + } + if v.Reconnects > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "reconnects", Code: "maximum"}) + } + if v.DrainTransitions != 0 && v.DrainTransitions < 0 { + violations = append(violations, FieldViolation{Field: "drain_transitions", Code: "minimum"}) + } + if v.DrainTransitions > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "drain_transitions", Code: "maximum"}) + } + if v.MediaDrops != 0 && v.MediaDrops < 0 { + violations = append(violations, FieldViolation{Field: "media_drops", Code: "minimum"}) + } + if v.MediaDrops > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "media_drops", Code: "maximum"}) + } + if v.MediaPackets != 0 && v.MediaPackets < 0 { + violations = append(violations, FieldViolation{Field: "media_packets", Code: "minimum"}) + } + if v.MediaPackets > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "media_packets", Code: "maximum"}) + } + if v.MediaBytes != 0 && v.MediaBytes < 0 { + violations = append(violations, FieldViolation{Field: "media_bytes", Code: "minimum"}) + } + if v.MediaBytes > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "media_bytes", Code: "maximum"}) + } + if v.QueueDelayMicros != 0 && v.QueueDelayMicros < 0 { + violations = append(violations, FieldViolation{Field: "queue_delay_micros", Code: "minimum"}) + } + if v.QueueDelayMicros > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "queue_delay_micros", Code: "maximum"}) + } + if v.ProcessingDelayMicros != 0 && v.ProcessingDelayMicros < 0 { + violations = append(violations, FieldViolation{Field: "processing_delay_micros", Code: "minimum"}) + } + if v.ProcessingDelayMicros > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "processing_delay_micros", Code: "maximum"}) + } + if v.ProcessingSamples != 0 && v.ProcessingSamples < 0 { + violations = append(violations, FieldViolation{Field: "processing_samples", Code: "minimum"}) + } + if v.ProcessingSamples > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "processing_samples", Code: "maximum"}) + } + if v.PacingDelayMicros != 0 && v.PacingDelayMicros < 0 { + violations = append(violations, FieldViolation{Field: "pacing_delay_micros", Code: "minimum"}) + } + if v.PacingDelayMicros > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "pacing_delay_micros", Code: "maximum"}) + } + if v.ProviderErrors != 0 && v.ProviderErrors < 0 { + violations = append(violations, FieldViolation{Field: "provider_errors", Code: "minimum"}) + } + if v.ProviderErrors > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "provider_errors", Code: "maximum"}) + } + if v.InputRejected != 0 && v.InputRejected < 0 { + violations = append(violations, FieldViolation{Field: "input_rejected", Code: "minimum"}) + } + if v.InputRejected > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "input_rejected", Code: "maximum"}) + } + if v.ControlRttMicros != 0 && v.ControlRttMicros < 0 { + violations = append(violations, FieldViolation{Field: "control_rtt_micros", Code: "minimum"}) + } + if v.ControlRttMicros > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "control_rtt_micros", Code: "maximum"}) + } + if v.ControlJitterMicros != 0 && v.ControlJitterMicros < 0 { + violations = append(violations, FieldViolation{Field: "control_jitter_micros", Code: "minimum"}) + } + if v.ControlJitterMicros > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "control_jitter_micros", Code: "maximum"}) + } + if v.ControlLossPpm != 0 && v.ControlLossPpm < 0 { + violations = append(violations, FieldViolation{Field: "control_loss_ppm", Code: "minimum"}) + } + if v.ControlLossPpm > 1000000 { + violations = append(violations, FieldViolation{Field: "control_loss_ppm", Code: "maximum"}) + } + if v.PendingReliable != 0 && v.PendingReliable < 0 { + violations = append(violations, FieldViolation{Field: "pending_reliable", Code: "minimum"}) + } + if v.PendingReliable > 9223372036854775807 { + violations = append(violations, FieldViolation{Field: "pending_reliable", Code: "maximum"}) + } + if v.ProviderState == "" { + violations = append(violations, FieldViolation{Field: "provider_state", Code: "required"}) + } + if v.ProviderState != "" && !(v.ProviderState == "unknown" || v.ProviderState == "starting" || v.ProviderState == "ready" || v.ProviderState == "disconnected" || v.ProviderState == "terminating" || v.ProviderState == "terminated" || v.ProviderState == "cleanup_pending" || v.ProviderState == "failed") { + violations = append(violations, FieldViolation{Field: "provider_state", Code: "invalid_value"}) + } + if len(violations) > 0 { + return ValidationError{Violations: violations} + } + return nil +} + +func DecodeGatewayTelemetry(data []byte) (GatewayTelemetry, error) { + var value GatewayTelemetry + if len(data) > 1024*1024 { + return value, errors.New("protocol payload exceeds limit") + } + var fields map[string]json.RawMessage + if err := json.Unmarshal(data, &fields); err != nil { + return value, err + } + if raw, ok := fields["admission_rejects"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "admission_rejects", Code: "required"}}} + } + if raw, ok := fields["admitted_sessions"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "admitted_sessions", Code: "required"}}} + } + if raw, ok := fields["control_jitter_micros"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "control_jitter_micros", Code: "required"}}} + } + if raw, ok := fields["control_loss_ppm"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "control_loss_ppm", Code: "required"}}} + } + if raw, ok := fields["control_rtt_micros"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "control_rtt_micros", Code: "required"}}} + } + if raw, ok := fields["drain_transitions"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "drain_transitions", Code: "required"}}} + } + if raw, ok := fields["input_rejected"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "input_rejected", Code: "required"}}} + } + if raw, ok := fields["media_bytes"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "media_bytes", Code: "required"}}} + } + if raw, ok := fields["media_drops"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "media_drops", Code: "required"}}} + } + if raw, ok := fields["media_packets"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "media_packets", Code: "required"}}} + } + if raw, ok := fields["pacing_delay_micros"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "pacing_delay_micros", Code: "required"}}} + } + if raw, ok := fields["pending_reliable"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "pending_reliable", Code: "required"}}} + } + if raw, ok := fields["processing_delay_micros"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "processing_delay_micros", Code: "required"}}} + } + if raw, ok := fields["processing_samples"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "processing_samples", Code: "required"}}} + } + if raw, ok := fields["provider_errors"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "provider_errors", Code: "required"}}} + } + if raw, ok := fields["provider_state"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "provider_state", Code: "required"}}} + } + if raw, ok := fields["queue_delay_micros"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "queue_delay_micros", Code: "required"}}} + } + if raw, ok := fields["reconnects"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "reconnects", Code: "required"}}} + } + decoder := json.NewDecoder(bytes.NewReader(data)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&value); err != nil { + return value, err + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + if err == nil { + return value, errors.New("trailing JSON value") + } + return value, err + } + if err := value.Validate(); err != nil { + return value, err + } + return value, nil +} + +func EncodeGatewayTelemetry(value GatewayTelemetry) ([]byte, error) { + if err := value.Validate(); err != nil { + return nil, err + } + return json.Marshal(value) +} + func (v GrantReference) Validate() error { var violations []FieldViolation if v.OpaqueValue == "" { @@ -3287,6 +3532,12 @@ func (v ProviderSessionWork) Validate() error { if len(v.PolicyVersionID) > 128 { violations = append(violations, FieldViolation{Field: "policy_version_id", Code: "max_length"}) } + if reflect.DeepEqual(v.StreamPolicy, ProviderStreamPolicy{}) { + violations = append(violations, FieldViolation{Field: "stream_policy", Code: "required"}) + } + if err := v.StreamPolicy.Validate(); err != nil { + violations = append(violations, FieldViolation{Field: "stream_policy", Code: "invalid_object"}) + } if v.ApplicationID == "" { violations = append(violations, FieldViolation{Field: "application_id", Code: "required"}) } @@ -3440,6 +3691,9 @@ func DecodeProviderSessionWork(data []byte) (ProviderSessionWork, error) { if raw, ok := fields["stream_host"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { return value, ValidationError{Violations: []FieldViolation{{Field: "stream_host", Code: "required"}}} } + if raw, ok := fields["stream_policy"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "stream_policy", Code: "required"}}} + } if raw, ok := fields["stream_port"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { return value, ValidationError{Violations: []FieldViolation{{Field: "stream_port", Code: "required"}}} } @@ -3555,6 +3809,108 @@ func EncodeProviderState(value ProviderState) ([]byte, error) { return json.Marshal(value) } +func (v ProviderStreamPolicy) Validate() error { + var violations []FieldViolation + if v.ResolutionWidth == 0 { + violations = append(violations, FieldViolation{Field: "resolution_width", Code: "required"}) + } + if v.ResolutionWidth != 0 && v.ResolutionWidth < 320 { + violations = append(violations, FieldViolation{Field: "resolution_width", Code: "minimum"}) + } + if v.ResolutionWidth > 16384 { + violations = append(violations, FieldViolation{Field: "resolution_width", Code: "maximum"}) + } + if v.ResolutionHeight == 0 { + violations = append(violations, FieldViolation{Field: "resolution_height", Code: "required"}) + } + if v.ResolutionHeight != 0 && v.ResolutionHeight < 200 { + violations = append(violations, FieldViolation{Field: "resolution_height", Code: "minimum"}) + } + if v.ResolutionHeight > 8640 { + violations = append(violations, FieldViolation{Field: "resolution_height", Code: "maximum"}) + } + if v.Fps == 0 { + violations = append(violations, FieldViolation{Field: "fps", Code: "required"}) + } + if v.Fps != 0 && v.Fps < 1 { + violations = append(violations, FieldViolation{Field: "fps", Code: "minimum"}) + } + if v.Fps > 240 { + violations = append(violations, FieldViolation{Field: "fps", Code: "maximum"}) + } + if v.Codec == "" { + violations = append(violations, FieldViolation{Field: "codec", Code: "required"}) + } + if v.Codec != "" && !(v.Codec == "H264" || v.Codec == "HEVC" || v.Codec == "AV1") { + violations = append(violations, FieldViolation{Field: "codec", Code: "invalid_value"}) + } + if v.BitrateKbps == 0 { + violations = append(violations, FieldViolation{Field: "bitrate_kbps", Code: "required"}) + } + if v.BitrateKbps != 0 && v.BitrateKbps < 100 { + violations = append(violations, FieldViolation{Field: "bitrate_kbps", Code: "minimum"}) + } + if v.BitrateKbps > 1000000 { + violations = append(violations, FieldViolation{Field: "bitrate_kbps", Code: "maximum"}) + } + if len(violations) > 0 { + return ValidationError{Violations: violations} + } + return nil +} + +func DecodeProviderStreamPolicy(data []byte) (ProviderStreamPolicy, error) { + var value ProviderStreamPolicy + if len(data) > 1024*1024 { + return value, errors.New("protocol payload exceeds limit") + } + var fields map[string]json.RawMessage + if err := json.Unmarshal(data, &fields); err != nil { + return value, err + } + if raw, ok := fields["audio_enabled"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "audio_enabled", Code: "required"}}} + } + if raw, ok := fields["bitrate_kbps"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "bitrate_kbps", Code: "required"}}} + } + if raw, ok := fields["codec"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "codec", Code: "required"}}} + } + if raw, ok := fields["fps"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "fps", Code: "required"}}} + } + if raw, ok := fields["resolution_height"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "resolution_height", Code: "required"}}} + } + if raw, ok := fields["resolution_width"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) { + return value, ValidationError{Violations: []FieldViolation{{Field: "resolution_width", Code: "required"}}} + } + decoder := json.NewDecoder(bytes.NewReader(data)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&value); err != nil { + return value, err + } + var trailing any + if err := decoder.Decode(&trailing); err != io.EOF { + if err == nil { + return value, errors.New("trailing JSON value") + } + return value, err + } + if err := value.Validate(); err != nil { + return value, err + } + return value, nil +} + +func EncodeProviderStreamPolicy(value ProviderStreamPolicy) ([]byte, error) { + if err := value.Validate(); err != nil { + return nil, err + } + return json.Marshal(value) +} + func (v ReauthGrant) Validate() error { var violations []FieldViolation if v.Token == "" { diff --git a/gen/manifest.json b/gen/manifest.json index 0c4e17c..21272b2 100644 --- a/gen/manifest.json +++ b/gen/manifest.json @@ -14,5 +14,5 @@ }, "generator_sha256": "922983e07a8ecc559771778fbf139155b14664742d9873be062880102777dccb", "protocol_version": "1.0.0", - "schema_sha256": "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" + "schema_sha256": "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b" } diff --git a/gen/rust/protocol.rs b/gen/rust/protocol.rs index 2a0e18f..3ed1942 100644 --- a/gen/rust/protocol.rs +++ b/gen/rust/protocol.rs @@ -1,6 +1,6 @@ // Code generated by tools/generate.py; DO NOT EDIT. #![allow(non_snake_case)] -pub const SCHEMA_SHA256: &str = "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8"; +pub const SCHEMA_SHA256: &str = "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b"; pub const CURRENT_WIRE_VERSION: &str = "1"; pub const N_MINUS_1_WIRE_VERSION: &str = "0"; pub const N_MINUS_2_WIRE_VERSION: &str = "-1"; @@ -771,11 +771,12 @@ pub struct GatewayHeartbeat { activeConnections: i64, egressKbps: i64, state: String, + telemetry: GatewayTelemetry, } impl GatewayHeartbeat { - pub fn new(version: String, gatewayId: String, sequence: i64, observedAt: String, activeConnections: i64, egressKbps: i64, state: String) -> Result { - let value = Self { version, gatewayId, sequence, observedAt, activeConnections, egressKbps, state }; + pub fn new(version: String, gatewayId: String, sequence: i64, observedAt: String, activeConnections: i64, egressKbps: i64, state: String, telemetry: GatewayTelemetry) -> Result { + let value = Self { version, gatewayId, sequence, observedAt, activeConnections, egressKbps, state, telemetry }; value.validate()?; Ok(value) } @@ -791,6 +792,7 @@ impl GatewayHeartbeat { if self.egressKbps < 0 { return Err(ValidationError::new("egress_kbps", "minimum")); } if self.egressKbps > 1000000000 { return Err(ValidationError::new("egress_kbps", "maximum")); } if self.state != "ready" && self.state != "draining" && self.state != "offline" { return Err(ValidationError::new("state", "invalid_value")); } + self.telemetry.validate().map_err(|_| ValidationError::new("telemetry", "invalid_object"))?; Ok(()) } pub fn version(&self) -> &String { &self.version } @@ -800,6 +802,7 @@ impl GatewayHeartbeat { pub fn activeConnections(&self) -> &i64 { &self.activeConnections } pub fn egressKbps(&self) -> &i64 { &self.egressKbps } pub fn state(&self) -> &String { &self.state } + pub fn telemetry(&self) -> &GatewayTelemetry { &self.telemetry } } #[derive(Debug, Clone, PartialEq, Eq)] @@ -873,6 +876,92 @@ impl GatewayRegistration { pub fn capabilities(&self) -> &CapabilityProfile { &self.capabilities } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct GatewayTelemetry { + admittedSessions: i64, + admissionRejects: i64, + reconnects: i64, + drainTransitions: i64, + mediaDrops: i64, + mediaPackets: i64, + mediaBytes: i64, + queueDelayMicros: i64, + processingDelayMicros: i64, + processingSamples: i64, + pacingDelayMicros: i64, + providerErrors: i64, + inputRejected: i64, + controlRttMicros: i64, + controlJitterMicros: i64, + controlLossPpm: i64, + pendingReliable: i64, + providerState: String, +} + +impl GatewayTelemetry { + pub fn new(admittedSessions: i64, admissionRejects: i64, reconnects: i64, drainTransitions: i64, mediaDrops: i64, mediaPackets: i64, mediaBytes: i64, queueDelayMicros: i64, processingDelayMicros: i64, processingSamples: i64, pacingDelayMicros: i64, providerErrors: i64, inputRejected: i64, controlRttMicros: i64, controlJitterMicros: i64, controlLossPpm: i64, pendingReliable: i64, providerState: String) -> Result { + let value = Self { admittedSessions, admissionRejects, reconnects, drainTransitions, mediaDrops, mediaPackets, mediaBytes, queueDelayMicros, processingDelayMicros, processingSamples, pacingDelayMicros, providerErrors, inputRejected, controlRttMicros, controlJitterMicros, controlLossPpm, pendingReliable, providerState }; + value.validate()?; + Ok(value) + } + pub fn validate(&self) -> Result<(), ValidationError> { + if self.admittedSessions < 0 { return Err(ValidationError::new("admitted_sessions", "minimum")); } + if self.admittedSessions > 9223372036854775807 { return Err(ValidationError::new("admitted_sessions", "maximum")); } + if self.admissionRejects < 0 { return Err(ValidationError::new("admission_rejects", "minimum")); } + if self.admissionRejects > 9223372036854775807 { return Err(ValidationError::new("admission_rejects", "maximum")); } + if self.reconnects < 0 { return Err(ValidationError::new("reconnects", "minimum")); } + if self.reconnects > 9223372036854775807 { return Err(ValidationError::new("reconnects", "maximum")); } + if self.drainTransitions < 0 { return Err(ValidationError::new("drain_transitions", "minimum")); } + if self.drainTransitions > 9223372036854775807 { return Err(ValidationError::new("drain_transitions", "maximum")); } + if self.mediaDrops < 0 { return Err(ValidationError::new("media_drops", "minimum")); } + if self.mediaDrops > 9223372036854775807 { return Err(ValidationError::new("media_drops", "maximum")); } + if self.mediaPackets < 0 { return Err(ValidationError::new("media_packets", "minimum")); } + if self.mediaPackets > 9223372036854775807 { return Err(ValidationError::new("media_packets", "maximum")); } + if self.mediaBytes < 0 { return Err(ValidationError::new("media_bytes", "minimum")); } + if self.mediaBytes > 9223372036854775807 { return Err(ValidationError::new("media_bytes", "maximum")); } + if self.queueDelayMicros < 0 { return Err(ValidationError::new("queue_delay_micros", "minimum")); } + if self.queueDelayMicros > 9223372036854775807 { return Err(ValidationError::new("queue_delay_micros", "maximum")); } + if self.processingDelayMicros < 0 { return Err(ValidationError::new("processing_delay_micros", "minimum")); } + if self.processingDelayMicros > 9223372036854775807 { return Err(ValidationError::new("processing_delay_micros", "maximum")); } + if self.processingSamples < 0 { return Err(ValidationError::new("processing_samples", "minimum")); } + if self.processingSamples > 9223372036854775807 { return Err(ValidationError::new("processing_samples", "maximum")); } + if self.pacingDelayMicros < 0 { return Err(ValidationError::new("pacing_delay_micros", "minimum")); } + if self.pacingDelayMicros > 9223372036854775807 { return Err(ValidationError::new("pacing_delay_micros", "maximum")); } + if self.providerErrors < 0 { return Err(ValidationError::new("provider_errors", "minimum")); } + if self.providerErrors > 9223372036854775807 { return Err(ValidationError::new("provider_errors", "maximum")); } + if self.inputRejected < 0 { return Err(ValidationError::new("input_rejected", "minimum")); } + if self.inputRejected > 9223372036854775807 { return Err(ValidationError::new("input_rejected", "maximum")); } + if self.controlRttMicros < 0 { return Err(ValidationError::new("control_rtt_micros", "minimum")); } + if self.controlRttMicros > 9223372036854775807 { return Err(ValidationError::new("control_rtt_micros", "maximum")); } + if self.controlJitterMicros < 0 { return Err(ValidationError::new("control_jitter_micros", "minimum")); } + if self.controlJitterMicros > 9223372036854775807 { return Err(ValidationError::new("control_jitter_micros", "maximum")); } + if self.controlLossPpm < 0 { return Err(ValidationError::new("control_loss_ppm", "minimum")); } + if self.controlLossPpm > 1000000 { return Err(ValidationError::new("control_loss_ppm", "maximum")); } + if self.pendingReliable < 0 { return Err(ValidationError::new("pending_reliable", "minimum")); } + if self.pendingReliable > 9223372036854775807 { return Err(ValidationError::new("pending_reliable", "maximum")); } + if self.providerState != "unknown" && self.providerState != "starting" && self.providerState != "ready" && self.providerState != "disconnected" && self.providerState != "terminating" && self.providerState != "terminated" && self.providerState != "cleanup_pending" && self.providerState != "failed" { return Err(ValidationError::new("provider_state", "invalid_value")); } + Ok(()) + } + pub fn admittedSessions(&self) -> &i64 { &self.admittedSessions } + pub fn admissionRejects(&self) -> &i64 { &self.admissionRejects } + pub fn reconnects(&self) -> &i64 { &self.reconnects } + pub fn drainTransitions(&self) -> &i64 { &self.drainTransitions } + pub fn mediaDrops(&self) -> &i64 { &self.mediaDrops } + pub fn mediaPackets(&self) -> &i64 { &self.mediaPackets } + pub fn mediaBytes(&self) -> &i64 { &self.mediaBytes } + pub fn queueDelayMicros(&self) -> &i64 { &self.queueDelayMicros } + pub fn processingDelayMicros(&self) -> &i64 { &self.processingDelayMicros } + pub fn processingSamples(&self) -> &i64 { &self.processingSamples } + pub fn pacingDelayMicros(&self) -> &i64 { &self.pacingDelayMicros } + pub fn providerErrors(&self) -> &i64 { &self.providerErrors } + pub fn inputRejected(&self) -> &i64 { &self.inputRejected } + pub fn controlRttMicros(&self) -> &i64 { &self.controlRttMicros } + pub fn controlJitterMicros(&self) -> &i64 { &self.controlJitterMicros } + pub fn controlLossPpm(&self) -> &i64 { &self.controlLossPpm } + pub fn pendingReliable(&self) -> &i64 { &self.pendingReliable } + pub fn providerState(&self) -> &String { &self.providerState } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct GrantReference { opaqueValue: String, @@ -1108,6 +1197,7 @@ pub struct ProviderSessionWork { providerProfile: String, providerIdentity: String, policyVersionId: String, + streamPolicy: ProviderStreamPolicy, applicationId: String, clientId: String, managementHost: String, @@ -1122,8 +1212,8 @@ pub struct ProviderSessionWork { } impl ProviderSessionWork { - pub fn new(version: String, sessionId: String, gatewayId: String, reconnectSequence: i64, expiresAt: String, providerProfile: String, providerIdentity: String, policyVersionId: String, applicationId: String, clientId: String, managementHost: String, managementPort: i64, streamHost: String, streamPort: i64, clientCertificatePem: String, clientPrivateKeyPem: String, serverCertificatePem: String, clipboardPolicy: ClipboardPolicy, providerApplicationTerminationAllowed: bool) -> Result { - let value = Self { version, sessionId, gatewayId, reconnectSequence, expiresAt, providerProfile, providerIdentity, policyVersionId, applicationId, clientId, managementHost, managementPort, streamHost, streamPort, clientCertificatePem, clientPrivateKeyPem, serverCertificatePem, clipboardPolicy, providerApplicationTerminationAllowed }; + pub fn new(version: String, sessionId: String, gatewayId: String, reconnectSequence: i64, expiresAt: String, providerProfile: String, providerIdentity: String, policyVersionId: String, streamPolicy: ProviderStreamPolicy, applicationId: String, clientId: String, managementHost: String, managementPort: i64, streamHost: String, streamPort: i64, clientCertificatePem: String, clientPrivateKeyPem: String, serverCertificatePem: String, clipboardPolicy: ClipboardPolicy, providerApplicationTerminationAllowed: bool) -> Result { + let value = Self { version, sessionId, gatewayId, reconnectSequence, expiresAt, providerProfile, providerIdentity, policyVersionId, streamPolicy, applicationId, clientId, managementHost, managementPort, streamHost, streamPort, clientCertificatePem, clientPrivateKeyPem, serverCertificatePem, clipboardPolicy, providerApplicationTerminationAllowed }; value.validate()?; Ok(value) } @@ -1144,6 +1234,7 @@ impl ProviderSessionWork { if self.policyVersionId.is_empty() { return Err(ValidationError::new("policy_version_id", "required")); } if !self.policyVersionId.is_empty() && self.policyVersionId.len() < 1 { return Err(ValidationError::new("policy_version_id", "min_length")); } if self.policyVersionId.len() > 128 { return Err(ValidationError::new("policy_version_id", "max_length")); } + self.streamPolicy.validate().map_err(|_| ValidationError::new("stream_policy", "invalid_object"))?; if self.applicationId.is_empty() { return Err(ValidationError::new("application_id", "required")); } if !self.applicationId.is_empty() && self.applicationId.len() < 1 { return Err(ValidationError::new("application_id", "min_length")); } if self.applicationId.len() > 128 { return Err(ValidationError::new("application_id", "max_length")); } @@ -1180,6 +1271,7 @@ impl ProviderSessionWork { pub fn providerProfile(&self) -> &String { &self.providerProfile } pub fn providerIdentity(&self) -> &String { &self.providerIdentity } pub fn policyVersionId(&self) -> &String { &self.policyVersionId } + pub fn streamPolicy(&self) -> &ProviderStreamPolicy { &self.streamPolicy } pub fn applicationId(&self) -> &String { &self.applicationId } pub fn clientId(&self) -> &String { &self.clientId } pub fn managementHost(&self) -> &String { &self.managementHost } @@ -1224,6 +1316,42 @@ impl ProviderState { pub fn channels(&self) -> &Vec { &self.channels } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ProviderStreamPolicy { + resolutionWidth: i64, + resolutionHeight: i64, + fps: i64, + codec: String, + bitrateKbps: i64, + audioEnabled: bool, +} + +impl ProviderStreamPolicy { + pub fn new(resolutionWidth: i64, resolutionHeight: i64, fps: i64, codec: String, bitrateKbps: i64, audioEnabled: bool) -> Result { + let value = Self { resolutionWidth, resolutionHeight, fps, codec, bitrateKbps, audioEnabled }; + value.validate()?; + Ok(value) + } + pub fn validate(&self) -> Result<(), ValidationError> { + if self.resolutionWidth < 320 { return Err(ValidationError::new("resolution_width", "minimum")); } + if self.resolutionWidth > 16384 { return Err(ValidationError::new("resolution_width", "maximum")); } + if self.resolutionHeight < 200 { return Err(ValidationError::new("resolution_height", "minimum")); } + if self.resolutionHeight > 8640 { return Err(ValidationError::new("resolution_height", "maximum")); } + if self.fps < 1 { return Err(ValidationError::new("fps", "minimum")); } + if self.fps > 240 { return Err(ValidationError::new("fps", "maximum")); } + if self.codec != "H264" && self.codec != "HEVC" && self.codec != "AV1" { return Err(ValidationError::new("codec", "invalid_value")); } + if self.bitrateKbps < 100 { return Err(ValidationError::new("bitrate_kbps", "minimum")); } + if self.bitrateKbps > 1000000 { return Err(ValidationError::new("bitrate_kbps", "maximum")); } + Ok(()) + } + pub fn resolutionWidth(&self) -> &i64 { &self.resolutionWidth } + pub fn resolutionHeight(&self) -> &i64 { &self.resolutionHeight } + pub fn fps(&self) -> &i64 { &self.fps } + pub fn codec(&self) -> &String { &self.codec } + pub fn bitrateKbps(&self) -> &i64 { &self.bitrateKbps } + pub fn audioEnabled(&self) -> &bool { &self.audioEnabled } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct ReauthGrant { token: String, diff --git a/gen/swift/Protocol.swift b/gen/swift/Protocol.swift index 32c28eb..f047c17 100644 --- a/gen/swift/Protocol.swift +++ b/gen/swift/Protocol.swift @@ -1,7 +1,7 @@ // Code generated by tools/generate.py; DO NOT EDIT. import Foundation public typealias JSONObject = [String: String] -public let schemaSHA256 = "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" +public let schemaSHA256 = "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b" public let currentWireVersion = "1" public let nMinus1WireVersion = "0" public let nMinus2WireVersion = "-1" @@ -1003,6 +1003,7 @@ public struct GatewayHeartbeat: Codable, Equatable { public let activeConnections: Int64 public let egressKbps: Int64 public let state: String + public let telemetry: GatewayTelemetry enum CodingKeys: String, CodingKey { case version = "version" case gatewayId = "gateway_id" @@ -1011,9 +1012,10 @@ public struct GatewayHeartbeat: Codable, Equatable { case activeConnections = "active_connections" case egressKbps = "egress_kbps" case state = "state" + case telemetry = "telemetry" } - public init(version: String, gatewayId: String, sequence: Int64, observedAt: String, activeConnections: Int64, egressKbps: Int64, state: String) throws { + public init(version: String, gatewayId: String, sequence: Int64, observedAt: String, activeConnections: Int64, egressKbps: Int64, state: String, telemetry: GatewayTelemetry) throws { self.version = version self.gatewayId = gatewayId self.sequence = sequence @@ -1021,6 +1023,7 @@ public struct GatewayHeartbeat: Codable, Equatable { self.activeConnections = activeConnections self.egressKbps = egressKbps self.state = state + self.telemetry = telemetry try validate() } @@ -1028,7 +1031,7 @@ public struct GatewayHeartbeat: Codable, Equatable { let all = try decoder.container(keyedBy: AnyCodingKey.self) for key in all.allKeys where CodingKeys(stringValue: key.stringValue) == nil { throw ContractValidationError(field: key.stringValue, code: "unknown_field") } let c = try decoder.container(keyedBy: CodingKeys.self) - try self.init(version: try c.decode(String.self, forKey: .version), gatewayId: try c.decode(String.self, forKey: .gatewayId), sequence: try c.decode(Int64.self, forKey: .sequence), observedAt: try c.decode(String.self, forKey: .observedAt), activeConnections: try c.decode(Int64.self, forKey: .activeConnections), egressKbps: try c.decode(Int64.self, forKey: .egressKbps), state: try c.decode(String.self, forKey: .state)) + try self.init(version: try c.decode(String.self, forKey: .version), gatewayId: try c.decode(String.self, forKey: .gatewayId), sequence: try c.decode(Int64.self, forKey: .sequence), observedAt: try c.decode(String.self, forKey: .observedAt), activeConnections: try c.decode(Int64.self, forKey: .activeConnections), egressKbps: try c.decode(Int64.self, forKey: .egressKbps), state: try c.decode(String.self, forKey: .state), telemetry: try c.decode(GatewayTelemetry.self, forKey: .telemetry)) } public func validate() throws { @@ -1044,6 +1047,7 @@ public struct GatewayHeartbeat: Codable, Equatable { if self.egressKbps < 0 { throw ContractValidationError(field: "egress_kbps", code: "minimum") } if self.egressKbps > 1000000000 { throw ContractValidationError(field: "egress_kbps", code: "maximum") } if !["ready", "draining", "offline"].contains(self.state) { throw ContractValidationError(field: "state", code: "invalid_value") } + try self.telemetry.validate() } public static func decodeJSON(_ data: Data) throws -> Self { try JSONDecoder().decode(Self.self, from: data) } @@ -1141,6 +1145,117 @@ public struct GatewayRegistration: Codable, Equatable { public func encodeJSON() throws -> Data { try validate(); return try JSONEncoder().encode(self) } } +public struct GatewayTelemetry: Codable, Equatable { + public let admittedSessions: Int64 + public let admissionRejects: Int64 + public let reconnects: Int64 + public let drainTransitions: Int64 + public let mediaDrops: Int64 + public let mediaPackets: Int64 + public let mediaBytes: Int64 + public let queueDelayMicros: Int64 + public let processingDelayMicros: Int64 + public let processingSamples: Int64 + public let pacingDelayMicros: Int64 + public let providerErrors: Int64 + public let inputRejected: Int64 + public let controlRttMicros: Int64 + public let controlJitterMicros: Int64 + public let controlLossPpm: Int64 + public let pendingReliable: Int64 + public let providerState: String + enum CodingKeys: String, CodingKey { + case admittedSessions = "admitted_sessions" + case admissionRejects = "admission_rejects" + case reconnects = "reconnects" + case drainTransitions = "drain_transitions" + case mediaDrops = "media_drops" + case mediaPackets = "media_packets" + case mediaBytes = "media_bytes" + case queueDelayMicros = "queue_delay_micros" + case processingDelayMicros = "processing_delay_micros" + case processingSamples = "processing_samples" + case pacingDelayMicros = "pacing_delay_micros" + case providerErrors = "provider_errors" + case inputRejected = "input_rejected" + case controlRttMicros = "control_rtt_micros" + case controlJitterMicros = "control_jitter_micros" + case controlLossPpm = "control_loss_ppm" + case pendingReliable = "pending_reliable" + case providerState = "provider_state" + } + + public init(admittedSessions: Int64, admissionRejects: Int64, reconnects: Int64, drainTransitions: Int64, mediaDrops: Int64, mediaPackets: Int64, mediaBytes: Int64, queueDelayMicros: Int64, processingDelayMicros: Int64, processingSamples: Int64, pacingDelayMicros: Int64, providerErrors: Int64, inputRejected: Int64, controlRttMicros: Int64, controlJitterMicros: Int64, controlLossPpm: Int64, pendingReliable: Int64, providerState: String) throws { + self.admittedSessions = admittedSessions + self.admissionRejects = admissionRejects + self.reconnects = reconnects + self.drainTransitions = drainTransitions + self.mediaDrops = mediaDrops + self.mediaPackets = mediaPackets + self.mediaBytes = mediaBytes + self.queueDelayMicros = queueDelayMicros + self.processingDelayMicros = processingDelayMicros + self.processingSamples = processingSamples + self.pacingDelayMicros = pacingDelayMicros + self.providerErrors = providerErrors + self.inputRejected = inputRejected + self.controlRttMicros = controlRttMicros + self.controlJitterMicros = controlJitterMicros + self.controlLossPpm = controlLossPpm + self.pendingReliable = pendingReliable + self.providerState = providerState + try validate() + } + + public init(from decoder: Decoder) throws { + let all = try decoder.container(keyedBy: AnyCodingKey.self) + for key in all.allKeys where CodingKeys(stringValue: key.stringValue) == nil { throw ContractValidationError(field: key.stringValue, code: "unknown_field") } + let c = try decoder.container(keyedBy: CodingKeys.self) + try self.init(admittedSessions: try c.decode(Int64.self, forKey: .admittedSessions), admissionRejects: try c.decode(Int64.self, forKey: .admissionRejects), reconnects: try c.decode(Int64.self, forKey: .reconnects), drainTransitions: try c.decode(Int64.self, forKey: .drainTransitions), mediaDrops: try c.decode(Int64.self, forKey: .mediaDrops), mediaPackets: try c.decode(Int64.self, forKey: .mediaPackets), mediaBytes: try c.decode(Int64.self, forKey: .mediaBytes), queueDelayMicros: try c.decode(Int64.self, forKey: .queueDelayMicros), processingDelayMicros: try c.decode(Int64.self, forKey: .processingDelayMicros), processingSamples: try c.decode(Int64.self, forKey: .processingSamples), pacingDelayMicros: try c.decode(Int64.self, forKey: .pacingDelayMicros), providerErrors: try c.decode(Int64.self, forKey: .providerErrors), inputRejected: try c.decode(Int64.self, forKey: .inputRejected), controlRttMicros: try c.decode(Int64.self, forKey: .controlRttMicros), controlJitterMicros: try c.decode(Int64.self, forKey: .controlJitterMicros), controlLossPpm: try c.decode(Int64.self, forKey: .controlLossPpm), pendingReliable: try c.decode(Int64.self, forKey: .pendingReliable), providerState: try c.decode(String.self, forKey: .providerState)) + } + + public func validate() throws { + if self.admittedSessions < 0 { throw ContractValidationError(field: "admitted_sessions", code: "minimum") } + if self.admittedSessions > 9223372036854775807 { throw ContractValidationError(field: "admitted_sessions", code: "maximum") } + if self.admissionRejects < 0 { throw ContractValidationError(field: "admission_rejects", code: "minimum") } + if self.admissionRejects > 9223372036854775807 { throw ContractValidationError(field: "admission_rejects", code: "maximum") } + if self.reconnects < 0 { throw ContractValidationError(field: "reconnects", code: "minimum") } + if self.reconnects > 9223372036854775807 { throw ContractValidationError(field: "reconnects", code: "maximum") } + if self.drainTransitions < 0 { throw ContractValidationError(field: "drain_transitions", code: "minimum") } + if self.drainTransitions > 9223372036854775807 { throw ContractValidationError(field: "drain_transitions", code: "maximum") } + if self.mediaDrops < 0 { throw ContractValidationError(field: "media_drops", code: "minimum") } + if self.mediaDrops > 9223372036854775807 { throw ContractValidationError(field: "media_drops", code: "maximum") } + if self.mediaPackets < 0 { throw ContractValidationError(field: "media_packets", code: "minimum") } + if self.mediaPackets > 9223372036854775807 { throw ContractValidationError(field: "media_packets", code: "maximum") } + if self.mediaBytes < 0 { throw ContractValidationError(field: "media_bytes", code: "minimum") } + if self.mediaBytes > 9223372036854775807 { throw ContractValidationError(field: "media_bytes", code: "maximum") } + if self.queueDelayMicros < 0 { throw ContractValidationError(field: "queue_delay_micros", code: "minimum") } + if self.queueDelayMicros > 9223372036854775807 { throw ContractValidationError(field: "queue_delay_micros", code: "maximum") } + if self.processingDelayMicros < 0 { throw ContractValidationError(field: "processing_delay_micros", code: "minimum") } + if self.processingDelayMicros > 9223372036854775807 { throw ContractValidationError(field: "processing_delay_micros", code: "maximum") } + if self.processingSamples < 0 { throw ContractValidationError(field: "processing_samples", code: "minimum") } + if self.processingSamples > 9223372036854775807 { throw ContractValidationError(field: "processing_samples", code: "maximum") } + if self.pacingDelayMicros < 0 { throw ContractValidationError(field: "pacing_delay_micros", code: "minimum") } + if self.pacingDelayMicros > 9223372036854775807 { throw ContractValidationError(field: "pacing_delay_micros", code: "maximum") } + if self.providerErrors < 0 { throw ContractValidationError(field: "provider_errors", code: "minimum") } + if self.providerErrors > 9223372036854775807 { throw ContractValidationError(field: "provider_errors", code: "maximum") } + if self.inputRejected < 0 { throw ContractValidationError(field: "input_rejected", code: "minimum") } + if self.inputRejected > 9223372036854775807 { throw ContractValidationError(field: "input_rejected", code: "maximum") } + if self.controlRttMicros < 0 { throw ContractValidationError(field: "control_rtt_micros", code: "minimum") } + if self.controlRttMicros > 9223372036854775807 { throw ContractValidationError(field: "control_rtt_micros", code: "maximum") } + if self.controlJitterMicros < 0 { throw ContractValidationError(field: "control_jitter_micros", code: "minimum") } + if self.controlJitterMicros > 9223372036854775807 { throw ContractValidationError(field: "control_jitter_micros", code: "maximum") } + if self.controlLossPpm < 0 { throw ContractValidationError(field: "control_loss_ppm", code: "minimum") } + if self.controlLossPpm > 1000000 { throw ContractValidationError(field: "control_loss_ppm", code: "maximum") } + if self.pendingReliable < 0 { throw ContractValidationError(field: "pending_reliable", code: "minimum") } + if self.pendingReliable > 9223372036854775807 { throw ContractValidationError(field: "pending_reliable", code: "maximum") } + if !["unknown", "starting", "ready", "disconnected", "terminating", "terminated", "cleanup_pending", "failed"].contains(self.providerState) { throw ContractValidationError(field: "provider_state", code: "invalid_value") } + } + + public static func decodeJSON(_ data: Data) throws -> Self { try JSONDecoder().decode(Self.self, from: data) } + public func encodeJSON() throws -> Data { try validate(); return try JSONEncoder().encode(self) } +} + public struct GrantReference: Codable, Equatable { public let opaqueValue: String public let expiresAt: String @@ -1458,6 +1573,7 @@ public struct ProviderSessionWork: Codable, Equatable { public let providerProfile: String public let providerIdentity: String public let policyVersionId: String + public let streamPolicy: ProviderStreamPolicy public let applicationId: String public let clientId: String public let managementHost: String @@ -1478,6 +1594,7 @@ public struct ProviderSessionWork: Codable, Equatable { case providerProfile = "provider_profile" case providerIdentity = "provider_identity" case policyVersionId = "policy_version_id" + case streamPolicy = "stream_policy" case applicationId = "application_id" case clientId = "client_id" case managementHost = "management_host" @@ -1491,7 +1608,7 @@ public struct ProviderSessionWork: Codable, Equatable { case providerApplicationTerminationAllowed = "provider_application_termination_allowed" } - public init(version: String, sessionId: String, gatewayId: String, reconnectSequence: Int64, expiresAt: String, providerProfile: String, providerIdentity: String, policyVersionId: String, applicationId: String, clientId: String, managementHost: String, managementPort: Int64, streamHost: String, streamPort: Int64, clientCertificatePem: String, clientPrivateKeyPem: String, serverCertificatePem: String, clipboardPolicy: ClipboardPolicy, providerApplicationTerminationAllowed: Bool) throws { + public init(version: String, sessionId: String, gatewayId: String, reconnectSequence: Int64, expiresAt: String, providerProfile: String, providerIdentity: String, policyVersionId: String, streamPolicy: ProviderStreamPolicy, applicationId: String, clientId: String, managementHost: String, managementPort: Int64, streamHost: String, streamPort: Int64, clientCertificatePem: String, clientPrivateKeyPem: String, serverCertificatePem: String, clipboardPolicy: ClipboardPolicy, providerApplicationTerminationAllowed: Bool) throws { self.version = version self.sessionId = sessionId self.gatewayId = gatewayId @@ -1500,6 +1617,7 @@ public struct ProviderSessionWork: Codable, Equatable { self.providerProfile = providerProfile self.providerIdentity = providerIdentity self.policyVersionId = policyVersionId + self.streamPolicy = streamPolicy self.applicationId = applicationId self.clientId = clientId self.managementHost = managementHost @@ -1518,7 +1636,7 @@ public struct ProviderSessionWork: Codable, Equatable { let all = try decoder.container(keyedBy: AnyCodingKey.self) for key in all.allKeys where CodingKeys(stringValue: key.stringValue) == nil { throw ContractValidationError(field: key.stringValue, code: "unknown_field") } let c = try decoder.container(keyedBy: CodingKeys.self) - try self.init(version: try c.decode(String.self, forKey: .version), sessionId: try c.decode(String.self, forKey: .sessionId), gatewayId: try c.decode(String.self, forKey: .gatewayId), reconnectSequence: try c.decode(Int64.self, forKey: .reconnectSequence), expiresAt: try c.decode(String.self, forKey: .expiresAt), providerProfile: try c.decode(String.self, forKey: .providerProfile), providerIdentity: try c.decode(String.self, forKey: .providerIdentity), policyVersionId: try c.decode(String.self, forKey: .policyVersionId), applicationId: try c.decode(String.self, forKey: .applicationId), clientId: try c.decode(String.self, forKey: .clientId), managementHost: try c.decode(String.self, forKey: .managementHost), managementPort: try c.decode(Int64.self, forKey: .managementPort), streamHost: try c.decode(String.self, forKey: .streamHost), streamPort: try c.decode(Int64.self, forKey: .streamPort), clientCertificatePem: try c.decode(String.self, forKey: .clientCertificatePem), clientPrivateKeyPem: try c.decode(String.self, forKey: .clientPrivateKeyPem), serverCertificatePem: try c.decode(String.self, forKey: .serverCertificatePem), clipboardPolicy: try c.decode(ClipboardPolicy.self, forKey: .clipboardPolicy), providerApplicationTerminationAllowed: try c.decode(Bool.self, forKey: .providerApplicationTerminationAllowed)) + try self.init(version: try c.decode(String.self, forKey: .version), sessionId: try c.decode(String.self, forKey: .sessionId), gatewayId: try c.decode(String.self, forKey: .gatewayId), reconnectSequence: try c.decode(Int64.self, forKey: .reconnectSequence), expiresAt: try c.decode(String.self, forKey: .expiresAt), providerProfile: try c.decode(String.self, forKey: .providerProfile), providerIdentity: try c.decode(String.self, forKey: .providerIdentity), policyVersionId: try c.decode(String.self, forKey: .policyVersionId), streamPolicy: try c.decode(ProviderStreamPolicy.self, forKey: .streamPolicy), applicationId: try c.decode(String.self, forKey: .applicationId), clientId: try c.decode(String.self, forKey: .clientId), managementHost: try c.decode(String.self, forKey: .managementHost), managementPort: try c.decode(Int64.self, forKey: .managementPort), streamHost: try c.decode(String.self, forKey: .streamHost), streamPort: try c.decode(Int64.self, forKey: .streamPort), clientCertificatePem: try c.decode(String.self, forKey: .clientCertificatePem), clientPrivateKeyPem: try c.decode(String.self, forKey: .clientPrivateKeyPem), serverCertificatePem: try c.decode(String.self, forKey: .serverCertificatePem), clipboardPolicy: try c.decode(ClipboardPolicy.self, forKey: .clipboardPolicy), providerApplicationTerminationAllowed: try c.decode(Bool.self, forKey: .providerApplicationTerminationAllowed)) } public func validate() throws { @@ -1539,6 +1657,7 @@ public struct ProviderSessionWork: Codable, Equatable { if self.policyVersionId.isEmpty { throw ContractValidationError(field: "policy_version_id", code: "required") } if !self.policyVersionId.isEmpty && self.policyVersionId.utf8.count < 1 { throw ContractValidationError(field: "policy_version_id", code: "min_length") } if self.policyVersionId.utf8.count > 128 { throw ContractValidationError(field: "policy_version_id", code: "max_length") } + try self.streamPolicy.validate() if self.applicationId.isEmpty { throw ContractValidationError(field: "application_id", code: "required") } if !self.applicationId.isEmpty && self.applicationId.utf8.count < 1 { throw ContractValidationError(field: "application_id", code: "min_length") } if self.applicationId.utf8.count > 128 { throw ContractValidationError(field: "application_id", code: "max_length") } @@ -1614,6 +1733,55 @@ public struct ProviderState: Codable, Equatable { public func encodeJSON() throws -> Data { try validate(); return try JSONEncoder().encode(self) } } +public struct ProviderStreamPolicy: Codable, Equatable { + public let resolutionWidth: Int64 + public let resolutionHeight: Int64 + public let fps: Int64 + public let codec: String + public let bitrateKbps: Int64 + public let audioEnabled: Bool + enum CodingKeys: String, CodingKey { + case resolutionWidth = "resolution_width" + case resolutionHeight = "resolution_height" + case fps = "fps" + case codec = "codec" + case bitrateKbps = "bitrate_kbps" + case audioEnabled = "audio_enabled" + } + + public init(resolutionWidth: Int64, resolutionHeight: Int64, fps: Int64, codec: String, bitrateKbps: Int64, audioEnabled: Bool) throws { + self.resolutionWidth = resolutionWidth + self.resolutionHeight = resolutionHeight + self.fps = fps + self.codec = codec + self.bitrateKbps = bitrateKbps + self.audioEnabled = audioEnabled + try validate() + } + + public init(from decoder: Decoder) throws { + let all = try decoder.container(keyedBy: AnyCodingKey.self) + for key in all.allKeys where CodingKeys(stringValue: key.stringValue) == nil { throw ContractValidationError(field: key.stringValue, code: "unknown_field") } + let c = try decoder.container(keyedBy: CodingKeys.self) + try self.init(resolutionWidth: try c.decode(Int64.self, forKey: .resolutionWidth), resolutionHeight: try c.decode(Int64.self, forKey: .resolutionHeight), fps: try c.decode(Int64.self, forKey: .fps), codec: try c.decode(String.self, forKey: .codec), bitrateKbps: try c.decode(Int64.self, forKey: .bitrateKbps), audioEnabled: try c.decode(Bool.self, forKey: .audioEnabled)) + } + + public func validate() throws { + if self.resolutionWidth < 320 { throw ContractValidationError(field: "resolution_width", code: "minimum") } + if self.resolutionWidth > 16384 { throw ContractValidationError(field: "resolution_width", code: "maximum") } + if self.resolutionHeight < 200 { throw ContractValidationError(field: "resolution_height", code: "minimum") } + if self.resolutionHeight > 8640 { throw ContractValidationError(field: "resolution_height", code: "maximum") } + if self.fps < 1 { throw ContractValidationError(field: "fps", code: "minimum") } + if self.fps > 240 { throw ContractValidationError(field: "fps", code: "maximum") } + if !["H264", "HEVC", "AV1"].contains(self.codec) { throw ContractValidationError(field: "codec", code: "invalid_value") } + if self.bitrateKbps < 100 { throw ContractValidationError(field: "bitrate_kbps", code: "minimum") } + if self.bitrateKbps > 1000000 { throw ContractValidationError(field: "bitrate_kbps", code: "maximum") } + } + + public static func decodeJSON(_ data: Data) throws -> Self { try JSONDecoder().decode(Self.self, from: data) } + public func encodeJSON() throws -> Data { try validate(); return try JSONEncoder().encode(self) } +} + public struct ReauthGrant: Codable, Equatable { public let token: String public let purpose: String diff --git a/openspec/changes/phase3c-gateway-heartbeat-telemetry/.openspec.yaml b/openspec/changes/phase3c-gateway-heartbeat-telemetry/.openspec.yaml new file mode 100644 index 0000000..f205fc7 --- /dev/null +++ b/openspec/changes/phase3c-gateway-heartbeat-telemetry/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-07-29 diff --git a/openspec/changes/phase3c-gateway-heartbeat-telemetry/design.md b/openspec/changes/phase3c-gateway-heartbeat-telemetry/design.md new file mode 100644 index 0000000..714d4b6 --- /dev/null +++ b/openspec/changes/phase3c-gateway-heartbeat-telemetry/design.md @@ -0,0 +1,36 @@ +## Context + +The Data Plane already collects process-wide atomic counters and gauges. Only active connections and egress Kbps cross the authenticated heartbeat boundary. + +## Goals / Non-Goals + +**Goals:** + +- Carry the existing low-cardinality observations with explicit units. +- Bound every numeric field and enumerate provider state. +- Keep registration capacity distinct from measured traffic. + +**Non-Goals:** + +- Add session, route, endpoint, credential, label, or payload fields. +- Define a new telemetry transport. +- Publish a Protocol version. + +## Decisions + +- Nest the values in required `GatewayTelemetry` so heartbeat telemetry is one strict atomic contract. +- Use cumulative counters and microsecond delay totals plus processing samples; consumers can derive rates/averages without losing raw observations. +- Keep loss as parts per million and provider state as a bounded enum. + +## Risks / Trade-offs + +- [Cumulative counters approach signed integer limits] → Bound at signed 64-bit and saturate consumer conversions. +- [New required object breaks RC6] → Test locally and publish only under separate immutable-version authorization. + +## Migration Plan + +Regenerate all bindings locally, update both consumers through the temporary workspace, and stop at the immutable publication boundary. + +## Open Questions + +None. diff --git a/openspec/changes/phase3c-gateway-heartbeat-telemetry/proposal.md b/openspec/changes/phase3c-gateway-heartbeat-telemetry/proposal.md new file mode 100644 index 0000000..7a8cc91 --- /dev/null +++ b/openspec/changes/phase3c-gateway-heartbeat-telemetry/proposal.md @@ -0,0 +1,23 @@ +## Why + +Gateway heartbeat currently carries active sessions and one egress rate but cannot transport the required observed process and provider-path telemetry. + +## What Changes + +- Add one required bounded low-cardinality telemetry object to authenticated gateway heartbeat. +- Cover counters, delay totals/samples, control RTT/loss/jitter, pending reliability, reconnects, and provider state. +- Keep configured capacity exclusively in registration and measured egress in heartbeat. + +## Capabilities + +### New Capabilities + +- `gateway-heartbeat-telemetry`: Authenticated heartbeats carry bounded observed gateway telemetry without routes, sessions, credentials, or payload data. + +### Modified Capabilities + +None. + +## Impact + +The control-v1 schema, generated Go/Rust/Swift bindings, conformance checks, and both unpublished consumers require coordinated local updates. RC6 remains unchanged. diff --git a/openspec/changes/phase3c-gateway-heartbeat-telemetry/specs/gateway-heartbeat-telemetry/spec.md b/openspec/changes/phase3c-gateway-heartbeat-telemetry/specs/gateway-heartbeat-telemetry/spec.md new file mode 100644 index 0000000..448ff02 --- /dev/null +++ b/openspec/changes/phase3c-gateway-heartbeat-telemetry/specs/gateway-heartbeat-telemetry/spec.md @@ -0,0 +1,15 @@ +## ADDED Requirements + +### Requirement: Heartbeat carries observed gateway telemetry +Every authenticated `GatewayHeartbeat` SHALL carry the bounded process-level counters, delay totals and samples, control RTT/loss/jitter, pending reliable work, reconnect count, and provider state defined by `GatewayTelemetry`. + +#### Scenario: Valid telemetry heartbeat +- **WHEN** a gateway reports its current observed snapshot +- **THEN** Go, Rust, and Swift bindings accept the same bounded low-cardinality values and units + +### Requirement: Heartbeat telemetry excludes sensitive dimensions +Heartbeat telemetry MUST reject unknown fields and MUST NOT include session, route, endpoint, credential, label, or payload values. + +#### Scenario: Secret or high-cardinality field is attempted +- **WHEN** a heartbeat contains an unregistered session, route, endpoint, credential, or payload field +- **THEN** strict contract validation rejects it before authenticated transport diff --git a/openspec/changes/phase3c-gateway-heartbeat-telemetry/tasks.md b/openspec/changes/phase3c-gateway-heartbeat-telemetry/tasks.md new file mode 100644 index 0000000..71e3f19 --- /dev/null +++ b/openspec/changes/phase3c-gateway-heartbeat-telemetry/tasks.md @@ -0,0 +1,11 @@ +## 1. Contract + +- [x] 1.1 Add bounded GatewayTelemetry to every heartbeat +- [x] 1.2 Add Go, Rust, and Swift strict conformance checks +- [x] 1.3 Regenerate bindings and prove deterministic output + +## 2. Consumer Boundary + +- [x] 2.1 Verify local Data Plane and Server consumers through a temporary workspace +- [ ] 2.2 Publish one new never-reused immutable Protocol version under separate authorization +- [ ] 2.3 Resolve from empty caches and pin exact checksums in both consumers diff --git a/openspec/changes/phase3c-provider-stream-policy/.openspec.yaml b/openspec/changes/phase3c-provider-stream-policy/.openspec.yaml new file mode 100644 index 0000000..f205fc7 --- /dev/null +++ b/openspec/changes/phase3c-provider-stream-policy/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-07-29 diff --git a/openspec/changes/phase3c-provider-stream-policy/design.md b/openspec/changes/phase3c-provider-stream-policy/design.md new file mode 100644 index 0000000..94875ee --- /dev/null +++ b/openspec/changes/phase3c-provider-stream-policy/design.md @@ -0,0 +1,37 @@ +## Context + +The Server resolves an immutable stream-policy version, but RC6 provider work carries only its identifier. The Data Plane consequently cannot distinguish the authorized settings from local defaults. + +## Goals / Non-Goals + +**Goals:** + +- Carry only the effective launch settings required by the provider boundary. +- Generate identical validation from the canonical schema for all bindings. +- Preserve the policy-version identifier for audit correlation. + +**Non-Goals:** + +- Publish or mutate RC6. +- Add provider-specific capability negotiation to Protocol. +- Expose provider work or policy internals to Verse clients. + +## Decisions + +- Use one required nested `ProviderStreamPolicy` value in `ProviderSessionWork`; this keeps the policy settings atomic and avoids repeating validation. +- Carry the Server-selected target bitrate rather than all policy bounds because Apollo ANNOUNCE consumes one configured bitrate. +- Permit canonical `H264`, `HEVC`, and `AV1` values in the contract. A provider implementation must reject values it cannot honor rather than silently downgrade them. +- Carry `audio_enabled` even though the current Apollo path cannot truthfully disable audio; the Data Plane must fail closed for that combination. + +## Risks / Trade-offs + +- [New required field breaks RC6 consumers] → Publish only under a separately authorized new immutable version and pin both consumers after empty-cache resolution. +- [Provider capabilities differ] → Validate the effective policy against the selected provider before readiness. + +## Migration Plan + +Regenerate and verify bindings locally, update both consumers through a temporary workspace only, then stop at the publication boundary. RC6 remains unchanged. + +## Open Questions + +None. diff --git a/openspec/changes/phase3c-provider-stream-policy/proposal.md b/openspec/changes/phase3c-provider-stream-policy/proposal.md new file mode 100644 index 0000000..036b685 --- /dev/null +++ b/openspec/changes/phase3c-provider-stream-policy/proposal.md @@ -0,0 +1,23 @@ +## Why + +Provider work identifies an immutable stream-policy version but omits the effective settings, allowing a gateway to launch Apollo with unrelated hard-coded media parameters. + +## What Changes + +- Add the effective resolution, frame rate, codec, selected bitrate, and audio policy to authenticated provider work. +- Require generated Go, Rust, and Swift bindings to validate the same bounded stream-policy contract. +- Keep the new contract unpublished until a new immutable Protocol version is separately authorized. + +## Capabilities + +### New Capabilities + +- `provider-stream-policy`: Authenticated provider work carries the exact effective stream policy consumed by the provider launch. + +### Modified Capabilities + +None. + +## Impact + +The control-v1 schema, generated bindings, conformance fixtures, and downstream Server and Data Plane consumers require coordinated local updates. RC6 remains immutable and unchanged. diff --git a/openspec/changes/phase3c-provider-stream-policy/specs/provider-stream-policy/spec.md b/openspec/changes/phase3c-provider-stream-policy/specs/provider-stream-policy/spec.md new file mode 100644 index 0000000..9e18773 --- /dev/null +++ b/openspec/changes/phase3c-provider-stream-policy/specs/provider-stream-policy/spec.md @@ -0,0 +1,15 @@ +## ADDED Requirements + +### Requirement: Provider work carries the effective stream policy +Authenticated `ProviderSessionWork` SHALL carry the immutable policy version and its effective resolution, frame rate, codec, target bitrate, and audio-enabled decision. + +#### Scenario: Gateway receives an effective policy +- **WHEN** the Server issues provider work for an admitted session +- **THEN** the work identifies the policy version and includes the effective bounded stream-policy values + +### Requirement: Stream-policy bindings share one strict contract +Generated Go, Rust, and Swift bindings MUST reject missing, unknown, out-of-range, or unsupported stream-policy wire values according to the canonical schema. + +#### Scenario: Invalid policy is rejected consistently +- **WHEN** provider work contains an unknown codec or a value outside the canonical bounds +- **THEN** every generated binding rejects the work before it can reach provider setup diff --git a/openspec/changes/phase3c-provider-stream-policy/tasks.md b/openspec/changes/phase3c-provider-stream-policy/tasks.md new file mode 100644 index 0000000..5afda4c --- /dev/null +++ b/openspec/changes/phase3c-provider-stream-policy/tasks.md @@ -0,0 +1,11 @@ +## 1. Contract + +- [x] 1.1 Add bounded effective stream policy to ProviderSessionWork +- [x] 1.2 Add Go, Rust, and Swift conformance coverage +- [x] 1.3 Regenerate bindings and prove deterministic output + +## 2. Consumer Boundary + +- [x] 2.1 Verify local Server and Data Plane consumers through a temporary workspace +- [ ] 2.2 Publish one new never-reused immutable Protocol version under separate authorization +- [ ] 2.3 Resolve from empty caches and pin exact checksums in both consumers diff --git a/schemas/control-v1.schema.json b/schemas/control-v1.schema.json index 3160256..cbbd62b 100644 --- a/schemas/control-v1.schema.json +++ b/schemas/control-v1.schema.json @@ -399,10 +399,35 @@ "capabilities": {"$ref": "#/$defs/CapabilityProfile"} } }, + "GatewayTelemetry": { + "type": "object", + "additionalProperties": false, + "required": ["admitted_sessions", "admission_rejects", "reconnects", "drain_transitions", "media_drops", "media_packets", "media_bytes", "queue_delay_micros", "processing_delay_micros", "processing_samples", "pacing_delay_micros", "provider_errors", "input_rejected", "control_rtt_micros", "control_jitter_micros", "control_loss_ppm", "pending_reliable", "provider_state"], + "properties": { + "admitted_sessions": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "admission_rejects": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "reconnects": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "drain_transitions": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "media_drops": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "media_packets": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "media_bytes": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "queue_delay_micros": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "processing_delay_micros": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "processing_samples": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "pacing_delay_micros": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "provider_errors": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "input_rejected": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "control_rtt_micros": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "control_jitter_micros": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "control_loss_ppm": {"type": "integer", "minimum": 0, "maximum": 1000000}, + "pending_reliable": {"type": "integer", "minimum": 0, "maximum": 9223372036854775807}, + "provider_state": {"type": "string", "enum": ["unknown", "starting", "ready", "disconnected", "terminating", "terminated", "cleanup_pending", "failed"]} + } + }, "GatewayHeartbeat": { "type": "object", "additionalProperties": false, - "required": ["version", "gateway_id", "sequence", "observed_at", "active_connections", "egress_kbps", "state"], + "required": ["version", "gateway_id", "sequence", "observed_at", "active_connections", "egress_kbps", "state", "telemetry"], "properties": { "version": {"type": "string", "const": "1"}, "gateway_id": {"type": "string", "minLength": 1, "maxLength": 128}, @@ -410,7 +435,8 @@ "observed_at": {"type": "string", "format": "date-time", "maxLength": 64}, "active_connections": {"type": "integer", "minimum": 0, "maximum": 1000000}, "egress_kbps": {"type": "integer", "minimum": 0, "maximum": 1000000000}, - "state": {"type": "string", "enum": ["ready", "draining", "offline"]} + "state": {"type": "string", "enum": ["ready", "draining", "offline"]}, + "telemetry": {"$ref": "#/$defs/GatewayTelemetry"} } }, "GatewayDrain": { @@ -457,10 +483,23 @@ "provider_identity": {"type": "string", "minLength": 1, "maxLength": 256} } }, + "ProviderStreamPolicy": { + "type": "object", + "additionalProperties": false, + "required": ["resolution_width", "resolution_height", "fps", "codec", "bitrate_kbps", "audio_enabled"], + "properties": { + "resolution_width": {"type": "integer", "minimum": 320, "maximum": 16384}, + "resolution_height": {"type": "integer", "minimum": 200, "maximum": 8640}, + "fps": {"type": "integer", "minimum": 1, "maximum": 240}, + "codec": {"type": "string", "enum": ["H264", "HEVC", "AV1"]}, + "bitrate_kbps": {"type": "integer", "minimum": 100, "maximum": 1000000}, + "audio_enabled": {"type": "boolean"} + } + }, "ProviderSessionWork": { "type": "object", "additionalProperties": false, - "required": ["version", "session_id", "gateway_id", "reconnect_sequence", "expires_at", "provider_profile", "provider_identity", "policy_version_id", "application_id", "client_id", "management_host", "management_port", "stream_host", "stream_port", "client_certificate_pem", "client_private_key_pem", "server_certificate_pem", "clipboard_policy", "provider_application_termination_allowed"], + "required": ["version", "session_id", "gateway_id", "reconnect_sequence", "expires_at", "provider_profile", "provider_identity", "policy_version_id", "stream_policy", "application_id", "client_id", "management_host", "management_port", "stream_host", "stream_port", "client_certificate_pem", "client_private_key_pem", "server_certificate_pem", "clipboard_policy", "provider_application_termination_allowed"], "properties": { "version": {"type": "string", "const": "1"}, "session_id": {"type": "string", "minLength": 1, "maxLength": 128}, @@ -470,6 +509,7 @@ "provider_profile": {"type": "string", "const": "apollo"}, "provider_identity": {"type": "string", "minLength": 1, "maxLength": 256}, "policy_version_id": {"type": "string", "minLength": 1, "maxLength": 128}, + "stream_policy": {"$ref": "#/$defs/ProviderStreamPolicy"}, "application_id": {"type": "string", "minLength": 1, "maxLength": 128}, "client_id": {"type": "string", "minLength": 1, "maxLength": 128}, "management_host": {"type": "string", "minLength": 1, "maxLength": 256}, diff --git a/tests/go/protocol_test.go b/tests/go/protocol_test.go index 1e40215..65a6957 100644 --- a/tests/go/protocol_test.go +++ b/tests/go/protocol_test.go @@ -62,6 +62,22 @@ func TestGatewayRegistrationRejectsInvertedProtocolBounds(t *testing.T) { } } +func TestGatewayHeartbeatCarriesBoundedObservedTelemetry(t *testing.T) { + valid := `{"version":"1","gateway_id":"gateway-1","sequence":1,"observed_at":"2099-01-01T00:00:00Z","active_connections":1,"egress_kbps":64,"state":"ready","telemetry":{"admitted_sessions":2,"admission_rejects":3,"reconnects":4,"drain_transitions":5,"media_drops":6,"media_packets":7,"media_bytes":8000,"queue_delay_micros":9,"processing_delay_micros":10,"processing_samples":11,"pacing_delay_micros":12,"provider_errors":13,"input_rejected":14,"control_rtt_micros":15,"control_jitter_micros":16,"control_loss_ppm":17,"pending_reliable":18,"provider_state":"ready"}}` + if _, err := protocol.DecodeGatewayHeartbeat([]byte(valid)); err != nil { + t.Fatalf("valid gateway heartbeat rejected: %v", err) + } + for _, invalid := range []string{ + strings.Replace(valid, `,"telemetry":{`, `,"session_id":"forbidden","telemetry":{`, 1), + strings.Replace(valid, `"control_loss_ppm":17`, `"control_loss_ppm":1000001`, 1), + strings.Replace(valid, `"provider_state":"ready"`, `"provider_state":"provider.example:47984"`, 1), + } { + if _, err := protocol.DecodeGatewayHeartbeat([]byte(invalid)); err == nil { + t.Fatalf("invalid gateway heartbeat accepted: %s", invalid) + } + } +} + func TestCapabilityIntersectionRejectsNoOverlap(t *testing.T) { first := protocol.CapabilityProfile{Transport: "quic-tls13", Framing: "datagram-v1", Media: "encoded", Audio: "encoded", SourceRateControl: "server", ClientDecode: "h264-opus"} if got, err := protocol.IntersectCapabilityProfiles(first, first); err != nil || got != first { @@ -114,7 +130,7 @@ func TestSessionAuthorityRejectsProviderRoute(t *testing.T) { } func TestProviderSessionWorkIsStrictAndSessionBound(t *testing.T) { - valid := `{"version":"1","session_id":"session-1","gateway_id":"gateway-1","reconnect_sequence":0,"expires_at":"2099-01-01T00:00:00Z","provider_profile":"apollo","provider_identity":"provider-1","policy_version_id":"policy-1","application_id":"42","client_id":"paired-client-1","management_host":"apollo.test","management_port":47990,"stream_host":"apollo.test","stream_port":47984,"client_certificate_pem":"certificate","client_private_key_pem":"private-key","server_certificate_pem":"server-certificate","clipboard_policy":{"client_to_provider_enabled":false,"provider_to_client_enabled":false,"max_text_bytes":65536,"max_updates_per_minute":30},"provider_application_termination_allowed":false}` + valid := `{"version":"1","session_id":"session-1","gateway_id":"gateway-1","reconnect_sequence":0,"expires_at":"2099-01-01T00:00:00Z","provider_profile":"apollo","provider_identity":"provider-1","policy_version_id":"policy-1","stream_policy":{"resolution_width":2560,"resolution_height":1440,"fps":120,"codec":"HEVC","bitrate_kbps":40000,"audio_enabled":true},"application_id":"42","client_id":"paired-client-1","management_host":"apollo.test","management_port":47990,"stream_host":"apollo.test","stream_port":47984,"client_certificate_pem":"certificate","client_private_key_pem":"private-key","server_certificate_pem":"server-certificate","clipboard_policy":{"client_to_provider_enabled":false,"provider_to_client_enabled":false,"max_text_bytes":65536,"max_updates_per_minute":30},"provider_application_termination_allowed":false}` if _, err := protocol.DecodeProviderSessionWork([]byte(valid)); err != nil { t.Fatalf("valid provider work rejected: %v", err) } @@ -124,6 +140,16 @@ func TestProviderSessionWorkIsStrictAndSessionBound(t *testing.T) { if _, err := protocol.DecodeProviderSessionWork([]byte(strings.Replace(valid, `,"clipboard_policy":{"client_to_provider_enabled":false,"provider_to_client_enabled":false,"max_text_bytes":65536,"max_updates_per_minute":30}`, "", 1))); err == nil { t.Fatal("provider work accepted missing clipboard policy") } + for _, invalid := range []string{ + strings.Replace(valid, `,"stream_policy":{"resolution_width":2560,"resolution_height":1440,"fps":120,"codec":"HEVC","bitrate_kbps":40000,"audio_enabled":true}`, "", 1), + strings.Replace(valid, `"fps":120`, `"fps":241`, 1), + strings.Replace(valid, `"codec":"HEVC"`, `"codec":"VP9"`, 1), + strings.Replace(valid, `"audio_enabled":true`, `"audio_enabled":true,"unknown":false`, 1), + } { + if _, err := protocol.DecodeProviderSessionWork([]byte(invalid)); err == nil { + t.Fatalf("provider work accepted invalid stream policy: %s", invalid) + } + } } func TestGatewayClipboardAuditIsMetadataOnlyAndStrict(t *testing.T) { diff --git a/tools/test_generated_contracts.py b/tools/test_generated_contracts.py index bde3840..c759114 100644 --- a/tools/test_generated_contracts.py +++ b/tools/test_generated_contracts.py @@ -85,6 +85,41 @@ do { ) fatalError("invalid allocation bounds were accepted") } catch { } +let streamPolicy = try ProviderStreamPolicy( + resolutionWidth: 2560, resolutionHeight: 1440, fps: 120, + codec: "HEVC", bitrateKbps: 40000, audioEnabled: true +) +guard streamPolicy.codec == "HEVC" else { fatalError("stream policy changed") } +for invalid in [ + { try ProviderStreamPolicy(resolutionWidth: 319, resolutionHeight: 1440, fps: 120, codec: "HEVC", bitrateKbps: 40000, audioEnabled: true) }, + { try ProviderStreamPolicy(resolutionWidth: 2560, resolutionHeight: 1440, fps: 241, codec: "HEVC", bitrateKbps: 40000, audioEnabled: true) }, + { try ProviderStreamPolicy(resolutionWidth: 2560, resolutionHeight: 1440, fps: 120, codec: "VP9", bitrateKbps: 40000, audioEnabled: true) }, +] { + do { + _ = try invalid() + fatalError("invalid stream policy was accepted") + } catch { } +} +let telemetry = try GatewayTelemetry( + admittedSessions: 1, admissionRejects: 2, reconnects: 3, drainTransitions: 4, + mediaDrops: 5, mediaPackets: 6, mediaBytes: 7, queueDelayMicros: 8, + processingDelayMicros: 9, processingSamples: 10, pacingDelayMicros: 11, + providerErrors: 12, inputRejected: 13, controlRttMicros: 14, + controlJitterMicros: 15, controlLossPpm: 16, pendingReliable: 17, + providerState: "ready" +) +guard telemetry.mediaBytes == 7 else { fatalError("gateway telemetry changed") } +do { + _ = try GatewayTelemetry( + admittedSessions: 1, admissionRejects: 2, reconnects: 3, drainTransitions: 4, + mediaDrops: 5, mediaPackets: 6, mediaBytes: 7, queueDelayMicros: 8, + processingDelayMicros: 9, processingSamples: 10, pacingDelayMicros: 11, + providerErrors: 12, inputRejected: 13, controlRttMicros: 14, + controlJitterMicros: 15, controlLossPpm: 1000001, pendingReliable: 17, + providerState: "ready" + ) + fatalError("invalid gateway telemetry was accepted") +} catch { } for text in [ String(repeating: "a", count: 65536), String(repeating: "é", count: 32768), @@ -149,6 +184,27 @@ fn main() { assert!(AllocationPolicy::new( 100, 50, 25, "standard".into(), "audience".into(), "verse".into(), 1, 60, 300, ).is_err()); + assert!(ProviderStreamPolicy::new( + 2560, 1440, 120, "HEVC".into(), 40000, true, + ).is_ok()); + assert!(ProviderStreamPolicy::new( + 319, 1440, 120, "HEVC".into(), 40000, true, + ).is_err()); + assert!(ProviderStreamPolicy::new( + 2560, 1440, 241, "HEVC".into(), 40000, true, + ).is_err()); + assert!(ProviderStreamPolicy::new( + 2560, 1440, 120, "VP9".into(), 40000, true, + ).is_err()); + assert!(GatewayTelemetry::new( + 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, "ready".into(), + ).is_ok()); + assert!(GatewayTelemetry::new( + 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 1000001, 17, "ready".into(), + ).is_err()); + assert!(GatewayTelemetry::new( + 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, "provider.example:47984".into(), + ).is_err()); for text in [ "a".repeat(65536), "é".repeat(32768),