Files
ziti/controller/network/inspect.go
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

329 lines
9.4 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 network
import (
"encoding/json"
"fmt"
"regexp"
"strings"
"sync"
"time"
"github.com/michaelquigley/pfxlog"
"github.com/openziti/channel/v5"
"github.com/openziti/channel/v5/protobufs"
"github.com/openziti/foundation/v2/concurrenz"
"github.com/openziti/foundation/v2/debugz"
"github.com/openziti/ziti/v2/common/inspect"
"github.com/openziti/ziti/v2/common/pb/ctrl_pb"
"github.com/openziti/ziti/v2/common/servermetrics"
"github.com/openziti/ziti/v2/controller/model"
"github.com/openziti/ziti/v2/controller/raft"
"github.com/openziti/ziti/v2/controller/xt"
)
type InspectResultValue struct {
AppId string
Name string
Value string
}
type InspectResult struct {
Success bool
Errors []string
Results []*InspectResultValue
}
func NewInspectionsManager(network *Network) *InspectionsManager {
return &InspectionsManager{
network: network,
}
}
type InspectionsManager struct {
network *Network
}
// InspectLocal runs inspect processing only on the local controller, without fanning out to
// routers or peer controllers.
func (self *InspectionsManager) InspectLocal(values []string) *InspectResult {
ctx := &inspectRequestContext{
network: self.network,
timeout: time.Second * 10,
requestedValues: values,
response: InspectResult{Success: true},
}
for _, requested := range values {
ctx.InspectLocal(requested)
}
return &ctx.response
}
func (self *InspectionsManager) Inspect(appRegex string, values []string) *InspectResult {
ctx := &inspectRequestContext{
network: self.network,
timeout: time.Second * 10,
requestedValues: values,
waitGroup: concurrenz.NewWaitGroup(),
appRegex: appRegex,
response: InspectResult{Success: true},
}
var err error
ctx.regex, err = regexp.Compile(appRegex)
if err != nil {
ctx.appendError(self.network.GetAppId(), err.Error())
return &ctx.response
}
return ctx.RunInspections()
}
type inspectRequestContext struct {
network *Network
timeout time.Duration
requestedValues []string
waitGroup concurrenz.WaitGroup
response InspectResult
appRegex string
regex *regexp.Regexp
lock sync.Mutex
complete bool
}
func (ctx *inspectRequestContext) RunInspections() *InspectResult {
log := pfxlog.Logger().
WithField("appRegex", ctx.appRegex).
WithField("values", ctx.requestedValues).
WithField("timeout", ctx.timeout)
ctx.inspectLocal()
for _, router := range ctx.network.AllConnectedRouters() {
ctx.inspectRouter(router)
}
for _, ch := range ctx.network.Dispatcher.GetPeers() {
ctx.inspectPeer(ch.Id(), ch)
}
if !ctx.waitGroup.WaitForDone(ctx.timeout) {
log.Info("inspect timed out, some values may be missing")
}
ctx.lock.Lock()
defer ctx.lock.Unlock()
ctx.complete = true
return &ctx.response
}
func (ctx *inspectRequestContext) inspectLocal() {
log := pfxlog.Logger().
WithField("appRegex", ctx.appRegex).
WithField("appId", ctx.network.GetAppId()).
WithField("values", ctx.requestedValues).
WithField("timeout", ctx.timeout)
if ctx.regex.MatchString(ctx.network.GetAppId()) {
log.Debug("inspect matched")
notifier := make(chan struct{})
ctx.waitGroup.AddNotifier(notifier)
go func() {
for _, requested := range ctx.requestedValues {
ctx.InspectLocal(requested)
}
close(notifier)
}()
} else {
log.Debug("inspect not matched")
}
}
func (ctx *inspectRequestContext) InspectLocal(name string) {
lc := strings.ToLower(name)
if lc == "stackdump" {
result := debugz.GenerateStack()
ctx.handleLocalStringResponse(name, &result, nil)
} else if strings.HasPrefix(lc, "metrics") {
msg := servermetrics.Poll(ctx.network.metricsRegistry)
ctx.handleLocalJsonResponse(name, msg)
} else if lc == "config" {
if rc, ok := ctx.network.config.(renderConfig); ok {
val, err := rc.RenderJsonConfig()
ctx.handleLocalStringResponse(name, &val, err)
}
} else if lc == "cluster-config" {
if src, ok := ctx.network.Dispatcher.(renderConfig); ok {
val, err := src.RenderJsonConfig()
ctx.handleLocalStringResponse(name, &val, err)
}
} else if lc == "connected-routers" {
var result []map[string]any
for _, r := range ctx.network.Router.AllConnected() {
status := map[string]any{}
status["id"] = r.Id
status["name"] = r.Name
status["version"] = r.VersionInfo.Version
status["connectTime"] = r.ConnectTime.Format(time.RFC3339)
status["underlays"] = r.Control.GetChannel().GetUnderlayCountsByType()
result = append(result, status)
}
ctx.handleLocalJsonResponse(name, result)
} else if lc == "connected-peers" {
if raftController, ok := ctx.network.Dispatcher.(*raft.Controller); ok {
members, err := raftController.ListMembers()
if err != nil {
ctx.appendError(ctx.network.GetAppId(), err.Error())
return
}
ctx.handleLocalJsonResponse(name, members)
}
} else if lc == inspect.PeerDialerKey {
if raftController, ok := ctx.network.Dispatcher.(*raft.Controller); ok {
matched, val, err := raftController.InspectPeerDialer(name)
if err != nil {
ctx.appendError(ctx.network.GetAppId(), err.Error())
} else if matched && val != nil {
ctx.appendValue(ctx.network.GetAppId(), name, *val)
}
}
} else if lc == "router-messaging" {
routerMessagingState, err := ctx.network.RouterMessaging.Inspect()
if err != nil {
ctx.appendError(ctx.network.GetAppId(), err.Error())
return
}
ctx.handleLocalJsonResponse(name, routerMessagingState)
} else if strings.HasPrefix(lc, "terminator-costs") {
state := &inspect.TerminatorCostDetails{}
xt.GlobalCosts().IterCosts(func(terminatorId string, cost xt.Cost) {
state.Terminators = append(state.Terminators, cost.Inspect(terminatorId))
})
ctx.handleLocalJsonResponse(name, state)
} else if lc == inspect.RouterIdentityConnectionStatusesKey {
result := ctx.network.env.GetManagers().Identity.GetConnectionTracker().Inspect()
ctx.handleLocalJsonResponse(name, result)
} else {
for _, inspectTarget := range ctx.network.inspectionTargets.Value() {
if handled, val, err := inspectTarget(lc); handled {
ctx.handleLocalStringResponse(name, val, err)
}
}
}
}
func (ctx *inspectRequestContext) handleLocalJsonResponse(key string, val interface{}) {
js, err := json.Marshal(val)
if err != nil {
ctx.appendError(ctx.network.GetAppId(), fmt.Errorf("failed to marshall %s to json (%w)", key, err).Error())
} else {
ctx.appendValue(ctx.network.GetAppId(), key, string(js))
}
}
func (ctx *inspectRequestContext) handleLocalStringResponse(key string, val *string, err error) {
if err != nil {
ctx.appendError(ctx.network.GetAppId(), err.Error())
} else if val != nil {
ctx.appendValue(ctx.network.GetAppId(), key, *val)
}
}
func (ctx *inspectRequestContext) inspectRouter(router *model.Router) {
log := pfxlog.Logger().
WithField("appRegex", ctx.appRegex).
WithField("routerId", router.Id).
WithField("values", ctx.requestedValues).
WithField("timeout", ctx.timeout)
if ctx.regex.MatchString(router.Id) || ctx.regex.MatchString(router.Name) {
log.Debug("inspect matched")
notifier := make(chan struct{})
ctx.waitGroup.AddNotifier(notifier)
go ctx.handleCtrlChanMessaging(router.Id, router.Control.GetLowPrioritySender(), notifier)
} else {
log.Debug("inspect not matched")
}
}
func (ctx *inspectRequestContext) inspectPeer(id string, ch channel.Channel) {
log := pfxlog.Logger().
WithField("appRegex", ctx.appRegex).
WithField("ctrlId", id).
WithField("values", ctx.requestedValues).
WithField("timeout", ctx.timeout)
if ctx.regex.MatchString(id) {
log.Debug("inspect matched")
notifier := make(chan struct{})
ctx.waitGroup.AddNotifier(notifier)
go ctx.handleCtrlChanMessaging(id, ch, notifier)
} else {
log.Debug("inspect not matched")
}
}
func (ctx *inspectRequestContext) handleCtrlChanMessaging(id string, ch channel.Sender, notifier chan struct{}) {
defer close(notifier)
request := &ctrl_pb.InspectRequest{RequestedValues: ctx.requestedValues}
resp := &ctrl_pb.InspectResponse{}
respMsg, err := protobufs.MarshalTyped(request).WithTimeout(ctx.timeout).SendForReply(ch)
err = protobufs.TypedResponse(resp).Unmarshall(respMsg, err)
if err != nil {
ctx.appendError(id, err.Error())
return
}
for _, err := range resp.Errors {
ctx.appendError(id, err)
}
for _, val := range resp.Values {
ctx.appendValue(id, val.Name, val.Value)
}
}
func (ctx *inspectRequestContext) appendValue(appId string, name string, value string) {
ctx.lock.Lock()
defer ctx.lock.Unlock()
if !ctx.complete {
ctx.response.Results = append(ctx.response.Results, &InspectResultValue{
AppId: appId,
Name: name,
Value: value,
})
}
}
func (ctx *inspectRequestContext) appendError(appId string, err string) {
ctx.lock.Lock()
defer ctx.lock.Unlock()
if !ctx.complete {
ctx.response.Success = false
ctx.response.Errors = append(ctx.response.Errors, fmt.Sprintf("%v: %v", appId, err))
}
}