Files
Paul Lorenz 187aa11f24 Own the metrics wire format in ziti. Fixes #4036
- adds a common/servermetrics package that owns the metrics MetricsMessage wire
  format and the reporting/usage subsystem (message builder, usage registry,
  interval and usage counters), wrapping the openziti/metrics Registry for
  metric collection
- moves the controllers metrics reporter into the router package and removes it
  from the shared metrics package, breaking a common -> router/env import cycle
- repoints controller and router consumers to common/servermetrics; base metric
  collection stays on openziti/metrics
- keeps the proto field numbers and the metrics content-type identical so the
  encoding is byte-compatible across the move, and uses a distinct proto package
  name so ziti's and the library's messages coexist without a global proto
  registry clash
- adds a round-trip test asserting wire compatibility with the library's
  MetricsMessage
- leaves openziti/metrics unchanged, so sdk-golang and the shared xgress data
  plane are unaffected
2026-06-29 22:37:25 -04:00

314 lines
8.6 KiB
Go

/*
Copyright NetFoundry Inc.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
https://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package ctrl_pb
import (
"fmt"
"strconv"
"strings"
"github.com/michaelquigley/pfxlog"
"github.com/openziti/channel/v5"
"github.com/openziti/ziti/v2/common/ctrl_msg"
"github.com/openziti/ziti/v2/common/servermetrics/metrics_pb"
"google.golang.org/protobuf/proto"
)
type Decoder struct{}
const DECODER = "ctrl"
func (d Decoder) Decode(msg *channel.Message) ([]byte, bool) {
switch msg.ContentType {
case int32(ContentType_CircuitRequestType):
circuitRequest := &CircuitRequest{}
if err := proto.Unmarshal(msg.Body, circuitRequest); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Circuit Request")
meta["ingressId"] = circuitRequest.IngressId
meta["service"] = circuitRequest.Service
headers := make([]string, 0)
for k := range circuitRequest.PeerData {
headers = append(headers, strconv.Itoa(int(k)))
}
meta["peerData"] = strings.Join(headers, ",")
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_CreateTerminatorRequestType):
createTerminator := &CreateTerminatorRequest{}
if err := proto.Unmarshal(msg.Body, createTerminator); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Create Terminator Request")
meta["terminator"] = terminatorToString(createTerminator)
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_RemoveTerminatorRequestType):
removeTerminator := &RemoveTerminatorRequest{}
if err := proto.Unmarshal(msg.Body, removeTerminator); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Remove Terminator Request")
meta["terminatorId"] = removeTerminator.TerminatorId
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_ValidateTerminatorsRequestType):
request := &ValidateTerminatorsRequest{}
if err := proto.Unmarshal(msg.Body, request); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Validate Terminators")
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_VerifyRouterType):
request := &VerifyRouter{}
if err := proto.Unmarshal(msg.Body, request); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Verify Router")
meta["routerId"] = request.RouterId
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ctrl_msg.CircuitSuccessType):
meta := channel.NewTraceMessageDecode(DECODER, "Circuit Success Response")
meta["circuitId"] = string(msg.Body)
meta["address"] = string(msg.Headers[ctrl_msg.CircuitSuccessAddressHeader])
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
case int32(ctrl_msg.CircuitFailedType):
meta := channel.NewTraceMessageDecode(DECODER, "Circuit Failed Response")
message := string(msg.Body)
if message != "" {
meta["message"] = message
}
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
case int32(ContentType_RouterLinksType):
links := &RouterLinks{}
if err := proto.Unmarshal(msg.Body, links); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "RouterLinks")
var linksD []map[string]interface{}
for _, link := range links.Links {
linksD = append(linksD, map[string]interface{}{
"id": link.Id,
"dest": link.DestRouterId,
})
}
meta["links"] = linksD
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_FaultType):
fault := &Fault{}
if err := proto.Unmarshal(msg.Body, fault); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Fault")
meta["subject"] = fault.Subject.String()
meta["id"] = fault.Id
meta["iteration"] = fault.Iteration
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_RouteType):
route := &Route{}
if err := proto.Unmarshal(msg.Body, route); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Route")
meta["circuitId"] = route.CircuitId
if route.Egress != nil {
meta["egress.address"] = route.Egress.Address
meta["egress.destination"] = route.Egress.Destination
}
for i, forward := range route.Forwards {
meta[fmt.Sprintf("forward[%d].srcAddress", i)] = forward.SrcAddress
meta[fmt.Sprintf("forward[%d].dstAddress", i)] = forward.DstAddress
meta[fmt.Sprintf("forward[%d].dstType", i)] = forward.DstType.String()
}
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
return nil, true
}
case int32(ContentType_UnrouteType):
unroute := &Unroute{}
if err := proto.Unmarshal(msg.Body, unroute); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Unroute")
meta["circuitId"] = unroute.CircuitId
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
}
case int32(ContentType_MetricsType):
metricsMsg := &metrics_pb.MetricsMessage{}
if err := proto.Unmarshal(msg.Body, metricsMsg); err == nil {
meta := channel.NewTraceMessageDecode(DECODER, "Metrics")
for name, metric := range metricsMsg.Histograms {
meta[name+".min"] = metric.Min
meta[name+".mean"] = metric.Mean
meta[name+".max"] = metric.Max
meta[name+".p95"] = metric.P95
meta[name+".p99"] = metric.P99
}
for name, metric := range metricsMsg.Meters {
meta[name+".count"] = metric.Count
meta[name+".meanRate"] = metric.MeanRate
meta[name+".m1Rate"] = metric.M1Rate
meta[name+".m5Rate"] = metric.M5Rate
meta[name+".m15Rate"] = metric.M15Rate
}
for name, counter := range metricsMsg.IntervalCounters {
for _, bucket := range counter.Buckets {
for key, val := range bucket.Values {
meta[name+"."+key+"["+strconv.FormatInt(bucket.IntervalStartUTC, 10)+"]"] = val
}
}
}
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
} else {
pfxlog.Logger().Errorf("unexpected error (%s)", err)
}
case int32(ctrl_msg.RouteResultType):
meta := channel.NewTraceMessageDecode(DECODER, "Route Result")
meta["circuitId"] = string(msg.Body)
meta["attempt"], _ = msg.GetUint32Header(ctrl_msg.RouteResultAttemptHeader)
success, _ := msg.GetBoolHeader(ctrl_msg.RouteResultSuccessHeader)
meta["success"] = success
if !success {
meta["errormsg"], _ = msg.GetStringHeader(ctrl_msg.RouteResultErrorHeader)
}
data, err := meta.MarshalTraceMessageDecode()
if err != nil {
return nil, true
}
return data, true
}
return nil, false
}
func terminatorToString(request *CreateTerminatorRequest) string {
return fmt.Sprintf("{serviceId=[%s], binding=[%s], address=[%v]}", request.ServiceId, request.Binding, request.Address)
}