feat(protocol): carry stream policy and gateway telemetry

This commit is contained in:
sechmachine
2026-07-30 01:45:48 +07:00
parent 59741761ce
commit d26f8b60f8
17 changed files with 991 additions and 42 deletions
+383 -27
View File
@@ -13,7 +13,7 @@ import (
"time" "time"
) )
const SchemaSHA256 = "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" const SchemaSHA256 = "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b"
const ProtocolVersion = "1.0.0" const ProtocolVersion = "1.0.0"
const CurrentWireVersion = "1" const CurrentWireVersion = "1"
const NMinus1WireVersion = "0" const NMinus1WireVersion = "0"
@@ -191,13 +191,14 @@ type GatewayDrain struct {
} }
type GatewayHeartbeat struct { type GatewayHeartbeat struct {
Version string `json:"version"` Version string `json:"version"`
GatewayID string `json:"gateway_id"` GatewayID string `json:"gateway_id"`
Sequence int64 `json:"sequence"` Sequence int64 `json:"sequence"`
ObservedAt string `json:"observed_at"` ObservedAt string `json:"observed_at"`
ActiveConnections int64 `json:"active_connections"` ActiveConnections int64 `json:"active_connections"`
EgressKbps int64 `json:"egress_kbps"` EgressKbps int64 `json:"egress_kbps"`
State string `json:"state"` State string `json:"state"`
Telemetry GatewayTelemetry `json:"telemetry"`
} }
type GatewayRegistration struct { type GatewayRegistration struct {
@@ -216,6 +217,27 @@ type GatewayRegistration struct {
Capabilities CapabilityProfile `json:"capabilities"` 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 { type GrantReference struct {
OpaqueValue string `json:"opaque_value"` OpaqueValue string `json:"opaque_value"`
ExpiresAt string `json:"expires_at"` ExpiresAt string `json:"expires_at"`
@@ -265,25 +287,26 @@ type PageInfo struct {
} }
type ProviderSessionWork struct { type ProviderSessionWork struct {
Version string `json:"version"` Version string `json:"version"`
SessionID string `json:"session_id"` SessionID string `json:"session_id"`
GatewayID string `json:"gateway_id"` GatewayID string `json:"gateway_id"`
ReconnectSequence int64 `json:"reconnect_sequence"` ReconnectSequence int64 `json:"reconnect_sequence"`
ExpiresAt string `json:"expires_at"` ExpiresAt string `json:"expires_at"`
ProviderProfile string `json:"provider_profile"` ProviderProfile string `json:"provider_profile"`
ProviderIdentity string `json:"provider_identity"` ProviderIdentity string `json:"provider_identity"`
PolicyVersionID string `json:"policy_version_id"` PolicyVersionID string `json:"policy_version_id"`
ApplicationID string `json:"application_id"` StreamPolicy ProviderStreamPolicy `json:"stream_policy"`
ClientID string `json:"client_id"` ApplicationID string `json:"application_id"`
ManagementHost string `json:"management_host"` ClientID string `json:"client_id"`
ManagementPort int64 `json:"management_port"` ManagementHost string `json:"management_host"`
StreamHost string `json:"stream_host"` ManagementPort int64 `json:"management_port"`
StreamPort int64 `json:"stream_port"` StreamHost string `json:"stream_host"`
ClientCertificatePem string `json:"client_certificate_pem"` StreamPort int64 `json:"stream_port"`
ClientPrivateKeyPem string `json:"client_private_key_pem"` ClientCertificatePem string `json:"client_certificate_pem"`
ServerCertificatePem string `json:"server_certificate_pem"` ClientPrivateKeyPem string `json:"client_private_key_pem"`
ClipboardPolicy ClipboardPolicy `json:"clipboard_policy"` ServerCertificatePem string `json:"server_certificate_pem"`
ProviderApplicationTerminationAllowed bool `json:"provider_application_termination_allowed"` ClipboardPolicy ClipboardPolicy `json:"clipboard_policy"`
ProviderApplicationTerminationAllowed bool `json:"provider_application_termination_allowed"`
} }
type ProviderState struct { type ProviderState struct {
@@ -294,6 +317,15 @@ type ProviderState struct {
Channels []string `json:"channels"` 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 { type ReauthGrant struct {
Token string `json:"token"` Token string `json:"token"`
Purpose string `json:"purpose"` 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") { if v.State != "" && !(v.State == "ready" || v.State == "draining" || v.State == "offline") {
violations = append(violations, FieldViolation{Field: "state", Code: "invalid_value"}) 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 { if len(violations) > 0 {
return ValidationError{Violations: violations} 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")) { if raw, ok := fields["state"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) {
return value, ValidationError{Violations: []FieldViolation{{Field: "state", Code: "required"}}} 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")) { if raw, ok := fields["version"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) {
return value, ValidationError{Violations: []FieldViolation{{Field: "version", Code: "required"}}} return value, ValidationError{Violations: []FieldViolation{{Field: "version", Code: "required"}}}
} }
@@ -2623,6 +2664,210 @@ func EncodeGatewayRegistration(value GatewayRegistration) ([]byte, error) {
return json.Marshal(value) 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 { func (v GrantReference) Validate() error {
var violations []FieldViolation var violations []FieldViolation
if v.OpaqueValue == "" { if v.OpaqueValue == "" {
@@ -3287,6 +3532,12 @@ func (v ProviderSessionWork) Validate() error {
if len(v.PolicyVersionID) > 128 { if len(v.PolicyVersionID) > 128 {
violations = append(violations, FieldViolation{Field: "policy_version_id", Code: "max_length"}) 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 == "" { if v.ApplicationID == "" {
violations = append(violations, FieldViolation{Field: "application_id", Code: "required"}) 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")) { if raw, ok := fields["stream_host"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) {
return value, ValidationError{Violations: []FieldViolation{{Field: "stream_host", Code: "required"}}} 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")) { if raw, ok := fields["stream_port"]; !ok || bytes.Equal(bytes.TrimSpace(raw), []byte("null")) {
return value, ValidationError{Violations: []FieldViolation{{Field: "stream_port", Code: "required"}}} return value, ValidationError{Violations: []FieldViolation{{Field: "stream_port", Code: "required"}}}
} }
@@ -3555,6 +3809,108 @@ func EncodeProviderState(value ProviderState) ([]byte, error) {
return json.Marshal(value) 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 { func (v ReauthGrant) Validate() error {
var violations []FieldViolation var violations []FieldViolation
if v.Token == "" { if v.Token == "" {
+1 -1
View File
@@ -14,5 +14,5 @@
}, },
"generator_sha256": "922983e07a8ecc559771778fbf139155b14664742d9873be062880102777dccb", "generator_sha256": "922983e07a8ecc559771778fbf139155b14664742d9873be062880102777dccb",
"protocol_version": "1.0.0", "protocol_version": "1.0.0",
"schema_sha256": "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" "schema_sha256": "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b"
} }
+133 -5
View File
@@ -1,6 +1,6 @@
// Code generated by tools/generate.py; DO NOT EDIT. // Code generated by tools/generate.py; DO NOT EDIT.
#![allow(non_snake_case)] #![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 CURRENT_WIRE_VERSION: &str = "1";
pub const N_MINUS_1_WIRE_VERSION: &str = "0"; pub const N_MINUS_1_WIRE_VERSION: &str = "0";
pub const N_MINUS_2_WIRE_VERSION: &str = "-1"; pub const N_MINUS_2_WIRE_VERSION: &str = "-1";
@@ -771,11 +771,12 @@ pub struct GatewayHeartbeat {
activeConnections: i64, activeConnections: i64,
egressKbps: i64, egressKbps: i64,
state: String, state: String,
telemetry: GatewayTelemetry,
} }
impl GatewayHeartbeat { impl GatewayHeartbeat {
pub fn new(version: String, gatewayId: String, sequence: i64, observedAt: String, activeConnections: i64, egressKbps: i64, state: String) -> Result<Self, ValidationError> { pub fn new(version: String, gatewayId: String, sequence: i64, observedAt: String, activeConnections: i64, egressKbps: i64, state: String, telemetry: GatewayTelemetry) -> Result<Self, ValidationError> {
let value = Self { version, gatewayId, sequence, observedAt, activeConnections, egressKbps, state }; let value = Self { version, gatewayId, sequence, observedAt, activeConnections, egressKbps, state, telemetry };
value.validate()?; value.validate()?;
Ok(value) Ok(value)
} }
@@ -791,6 +792,7 @@ impl GatewayHeartbeat {
if self.egressKbps < 0 { return Err(ValidationError::new("egress_kbps", "minimum")); } if self.egressKbps < 0 { return Err(ValidationError::new("egress_kbps", "minimum")); }
if self.egressKbps > 1000000000 { return Err(ValidationError::new("egress_kbps", "maximum")); } 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")); } 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(()) Ok(())
} }
pub fn version(&self) -> &String { &self.version } pub fn version(&self) -> &String { &self.version }
@@ -800,6 +802,7 @@ impl GatewayHeartbeat {
pub fn activeConnections(&self) -> &i64 { &self.activeConnections } pub fn activeConnections(&self) -> &i64 { &self.activeConnections }
pub fn egressKbps(&self) -> &i64 { &self.egressKbps } pub fn egressKbps(&self) -> &i64 { &self.egressKbps }
pub fn state(&self) -> &String { &self.state } pub fn state(&self) -> &String { &self.state }
pub fn telemetry(&self) -> &GatewayTelemetry { &self.telemetry }
} }
#[derive(Debug, Clone, PartialEq, Eq)] #[derive(Debug, Clone, PartialEq, Eq)]
@@ -873,6 +876,92 @@ impl GatewayRegistration {
pub fn capabilities(&self) -> &CapabilityProfile { &self.capabilities } 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<Self, ValidationError> {
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)] #[derive(Debug, Clone, PartialEq, Eq)]
pub struct GrantReference { pub struct GrantReference {
opaqueValue: String, opaqueValue: String,
@@ -1108,6 +1197,7 @@ pub struct ProviderSessionWork {
providerProfile: String, providerProfile: String,
providerIdentity: String, providerIdentity: String,
policyVersionId: String, policyVersionId: String,
streamPolicy: ProviderStreamPolicy,
applicationId: String, applicationId: String,
clientId: String, clientId: String,
managementHost: String, managementHost: String,
@@ -1122,8 +1212,8 @@ pub struct ProviderSessionWork {
} }
impl 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<Self, ValidationError> { 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<Self, ValidationError> {
let value = Self { version, sessionId, gatewayId, reconnectSequence, expiresAt, providerProfile, providerIdentity, policyVersionId, applicationId, clientId, managementHost, managementPort, streamHost, streamPort, clientCertificatePem, clientPrivateKeyPem, serverCertificatePem, clipboardPolicy, providerApplicationTerminationAllowed }; 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()?; value.validate()?;
Ok(value) 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() { 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.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")); } 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() { 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.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")); } 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 providerProfile(&self) -> &String { &self.providerProfile }
pub fn providerIdentity(&self) -> &String { &self.providerIdentity } pub fn providerIdentity(&self) -> &String { &self.providerIdentity }
pub fn policyVersionId(&self) -> &String { &self.policyVersionId } pub fn policyVersionId(&self) -> &String { &self.policyVersionId }
pub fn streamPolicy(&self) -> &ProviderStreamPolicy { &self.streamPolicy }
pub fn applicationId(&self) -> &String { &self.applicationId } pub fn applicationId(&self) -> &String { &self.applicationId }
pub fn clientId(&self) -> &String { &self.clientId } pub fn clientId(&self) -> &String { &self.clientId }
pub fn managementHost(&self) -> &String { &self.managementHost } pub fn managementHost(&self) -> &String { &self.managementHost }
@@ -1224,6 +1316,42 @@ impl ProviderState {
pub fn channels(&self) -> &Vec<String> { &self.channels } pub fn channels(&self) -> &Vec<String> { &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<Self, ValidationError> {
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)] #[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReauthGrant { pub struct ReauthGrant {
token: String, token: String,
+173 -5
View File
@@ -1,7 +1,7 @@
// Code generated by tools/generate.py; DO NOT EDIT. // Code generated by tools/generate.py; DO NOT EDIT.
import Foundation import Foundation
public typealias JSONObject = [String: String] public typealias JSONObject = [String: String]
public let schemaSHA256 = "e98c75ef81bbeac6be2b8f11202c1ffecec0aa515b48576a26756290e99d5dd8" public let schemaSHA256 = "fe4d6be69665f09f04e2cd89273ebc5d2dddf4cde98b94a9d8c4f726504ffe1b"
public let currentWireVersion = "1" public let currentWireVersion = "1"
public let nMinus1WireVersion = "0" public let nMinus1WireVersion = "0"
public let nMinus2WireVersion = "-1" public let nMinus2WireVersion = "-1"
@@ -1003,6 +1003,7 @@ public struct GatewayHeartbeat: Codable, Equatable {
public let activeConnections: Int64 public let activeConnections: Int64
public let egressKbps: Int64 public let egressKbps: Int64
public let state: String public let state: String
public let telemetry: GatewayTelemetry
enum CodingKeys: String, CodingKey { enum CodingKeys: String, CodingKey {
case version = "version" case version = "version"
case gatewayId = "gateway_id" case gatewayId = "gateway_id"
@@ -1011,9 +1012,10 @@ public struct GatewayHeartbeat: Codable, Equatable {
case activeConnections = "active_connections" case activeConnections = "active_connections"
case egressKbps = "egress_kbps" case egressKbps = "egress_kbps"
case state = "state" 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.version = version
self.gatewayId = gatewayId self.gatewayId = gatewayId
self.sequence = sequence self.sequence = sequence
@@ -1021,6 +1023,7 @@ public struct GatewayHeartbeat: Codable, Equatable {
self.activeConnections = activeConnections self.activeConnections = activeConnections
self.egressKbps = egressKbps self.egressKbps = egressKbps
self.state = state self.state = state
self.telemetry = telemetry
try validate() try validate()
} }
@@ -1028,7 +1031,7 @@ public struct GatewayHeartbeat: Codable, Equatable {
let all = try decoder.container(keyedBy: AnyCodingKey.self) 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") } 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) 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 { 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 < 0 { throw ContractValidationError(field: "egress_kbps", code: "minimum") }
if self.egressKbps > 1000000000 { throw ContractValidationError(field: "egress_kbps", code: "maximum") } 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") } 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) } 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 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 struct GrantReference: Codable, Equatable {
public let opaqueValue: String public let opaqueValue: String
public let expiresAt: String public let expiresAt: String
@@ -1458,6 +1573,7 @@ public struct ProviderSessionWork: Codable, Equatable {
public let providerProfile: String public let providerProfile: String
public let providerIdentity: String public let providerIdentity: String
public let policyVersionId: String public let policyVersionId: String
public let streamPolicy: ProviderStreamPolicy
public let applicationId: String public let applicationId: String
public let clientId: String public let clientId: String
public let managementHost: String public let managementHost: String
@@ -1478,6 +1594,7 @@ public struct ProviderSessionWork: Codable, Equatable {
case providerProfile = "provider_profile" case providerProfile = "provider_profile"
case providerIdentity = "provider_identity" case providerIdentity = "provider_identity"
case policyVersionId = "policy_version_id" case policyVersionId = "policy_version_id"
case streamPolicy = "stream_policy"
case applicationId = "application_id" case applicationId = "application_id"
case clientId = "client_id" case clientId = "client_id"
case managementHost = "management_host" case managementHost = "management_host"
@@ -1491,7 +1608,7 @@ public struct ProviderSessionWork: Codable, Equatable {
case providerApplicationTerminationAllowed = "provider_application_termination_allowed" 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.version = version
self.sessionId = sessionId self.sessionId = sessionId
self.gatewayId = gatewayId self.gatewayId = gatewayId
@@ -1500,6 +1617,7 @@ public struct ProviderSessionWork: Codable, Equatable {
self.providerProfile = providerProfile self.providerProfile = providerProfile
self.providerIdentity = providerIdentity self.providerIdentity = providerIdentity
self.policyVersionId = policyVersionId self.policyVersionId = policyVersionId
self.streamPolicy = streamPolicy
self.applicationId = applicationId self.applicationId = applicationId
self.clientId = clientId self.clientId = clientId
self.managementHost = managementHost self.managementHost = managementHost
@@ -1518,7 +1636,7 @@ public struct ProviderSessionWork: Codable, Equatable {
let all = try decoder.container(keyedBy: AnyCodingKey.self) 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") } 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) 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 { 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 { 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.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") } 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 { 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.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") } 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 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 struct ReauthGrant: Codable, Equatable {
public let token: String public let token: String
public let purpose: String public let purpose: String
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-07-29
@@ -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.
@@ -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.
@@ -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
@@ -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
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-07-29
@@ -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.
@@ -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.
@@ -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
@@ -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
+43 -3
View File
@@ -399,10 +399,35 @@
"capabilities": {"$ref": "#/$defs/CapabilityProfile"} "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": { "GatewayHeartbeat": {
"type": "object", "type": "object",
"additionalProperties": false, "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": { "properties": {
"version": {"type": "string", "const": "1"}, "version": {"type": "string", "const": "1"},
"gateway_id": {"type": "string", "minLength": 1, "maxLength": 128}, "gateway_id": {"type": "string", "minLength": 1, "maxLength": 128},
@@ -410,7 +435,8 @@
"observed_at": {"type": "string", "format": "date-time", "maxLength": 64}, "observed_at": {"type": "string", "format": "date-time", "maxLength": 64},
"active_connections": {"type": "integer", "minimum": 0, "maximum": 1000000}, "active_connections": {"type": "integer", "minimum": 0, "maximum": 1000000},
"egress_kbps": {"type": "integer", "minimum": 0, "maximum": 1000000000}, "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": { "GatewayDrain": {
@@ -457,10 +483,23 @@
"provider_identity": {"type": "string", "minLength": 1, "maxLength": 256} "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": { "ProviderSessionWork": {
"type": "object", "type": "object",
"additionalProperties": false, "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": { "properties": {
"version": {"type": "string", "const": "1"}, "version": {"type": "string", "const": "1"},
"session_id": {"type": "string", "minLength": 1, "maxLength": 128}, "session_id": {"type": "string", "minLength": 1, "maxLength": 128},
@@ -470,6 +509,7 @@
"provider_profile": {"type": "string", "const": "apollo"}, "provider_profile": {"type": "string", "const": "apollo"},
"provider_identity": {"type": "string", "minLength": 1, "maxLength": 256}, "provider_identity": {"type": "string", "minLength": 1, "maxLength": 256},
"policy_version_id": {"type": "string", "minLength": 1, "maxLength": 128}, "policy_version_id": {"type": "string", "minLength": 1, "maxLength": 128},
"stream_policy": {"$ref": "#/$defs/ProviderStreamPolicy"},
"application_id": {"type": "string", "minLength": 1, "maxLength": 128}, "application_id": {"type": "string", "minLength": 1, "maxLength": 128},
"client_id": {"type": "string", "minLength": 1, "maxLength": 128}, "client_id": {"type": "string", "minLength": 1, "maxLength": 128},
"management_host": {"type": "string", "minLength": 1, "maxLength": 256}, "management_host": {"type": "string", "minLength": 1, "maxLength": 256},
+27 -1
View File
@@ -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) { func TestCapabilityIntersectionRejectsNoOverlap(t *testing.T) {
first := protocol.CapabilityProfile{Transport: "quic-tls13", Framing: "datagram-v1", Media: "encoded", Audio: "encoded", SourceRateControl: "server", ClientDecode: "h264-opus"} 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 { 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) { 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 { if _, err := protocol.DecodeProviderSessionWork([]byte(valid)); err != nil {
t.Fatalf("valid provider work rejected: %v", err) 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 { 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") 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) { func TestGatewayClipboardAuditIsMetadataOnlyAndStrict(t *testing.T) {
+56
View File
@@ -85,6 +85,41 @@ do {
) )
fatalError("invalid allocation bounds were accepted") fatalError("invalid allocation bounds were accepted")
} catch { } } 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 [ for text in [
String(repeating: "a", count: 65536), String(repeating: "a", count: 65536),
String(repeating: "é", count: 32768), String(repeating: "é", count: 32768),
@@ -149,6 +184,27 @@ fn main() {
assert!(AllocationPolicy::new( assert!(AllocationPolicy::new(
100, 50, 25, "standard".into(), "audience".into(), "verse".into(), 1, 60, 300, 100, 50, 25, "standard".into(), "audience".into(), "verse".into(), 1, 60, 300,
).is_err()); ).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 [ for text in [
"a".repeat(65536), "a".repeat(65536),
"é".repeat(32768), "é".repeat(32768),