diff --git a/controller/api/context.go b/controller/api/context.go index d8967a051..1b129a365 100644 --- a/controller/api/context.go +++ b/controller/api/context.go @@ -79,11 +79,22 @@ func (rc *RequestContextImpl) GetEntitySubId() (string, error) { } func (rc *RequestContextImpl) NewChangeContext() *change.Context { - src := fmt.Sprintf("rest[auth=fabric/host=%v/method=%v/remote=%v]", rc.GetRequest().Host, rc.GetRequest().Method, rc.GetRequest().RemoteAddr) - changeCtx := change.New().SetSource(src).SetChangeAuthorType("fabric.admin") - if rc.Request.Form.Has("traceId") { - changeCtx.SetChangeAuthorId(rc.Request.Form.Get("traceId")) + changeCtx := change.New().SetSourceType(change.SourceTypeRest). + SetSourceAuth("fabric"). + SetSourceMethod(rc.GetRequest().Method). + SetSourceLocal(rc.GetRequest().Host). + SetSourceRemote(rc.GetRequest().RemoteAddr) + + changeCtx.SetChangeAuthorType(change.AuthorTypeCert) + + if rc.Request.TLS != nil { + for _, cert := range rc.Request.TLS.PeerCertificates { + if !cert.IsCA { + changeCtx.SetChangeAuthorId(cert.Subject.CommonName) + } + } } + return changeCtx } diff --git a/controller/change/context.go b/controller/change/context.go index e9acfec09..b91d5588a 100644 --- a/controller/change/context.go +++ b/controller/change/context.go @@ -31,7 +31,27 @@ const ( AuthorNameKey = "authorName" AuthorTypeKey = "authorType" TraceIdKey = "traceId" - Source = "source" + SourceType = "src.type" + SourceAuth = "src.auth" + SourceMethod = "src.method" + SourceLocal = "src.local" + SourceRemote = "src.remote" +) + +type AuthorType string + +const ( + AuthorTypeCert = "cert" + AuthorTypeIdentity = "identity" + AuthorTypeRouter = "router" + AuthorTypeController = "controller" + AuthorTypeUnattributed = "unattributed" +) + +const ( + SourceTypeControlChannel = "ctrl.channel" + SourceTypeRest = "rest" + SourceTypeXt = "xt" ) func New() *Context { @@ -45,6 +65,25 @@ type Context struct { RaftIndex uint64 } +type Author struct { + Type string `json:"type"` + Id string `json:"id,omitempty"` + Name string `json:"name,omitempty"` +} + +type Source struct { + Type string `json:"type"` + Auth string `json:"auth,omitempty"` + LocalAddr string `json:"local_addr,omitempty"` + RemoteAddr string `json:"remote_addr,omitempty"` + Method string `json:"method,omitempty"` +} + +func (self *Context) SetChangeAuthorType(val AuthorType) *Context { + self.Attributes[AuthorTypeKey] = string(val) + return self +} + func (self *Context) SetChangeAuthorId(val string) *Context { self.Attributes[AuthorIdKey] = val return self @@ -55,25 +94,73 @@ func (self *Context) SetChangeAuthorName(val string) *Context { return self } -func (self *Context) SetChangeAuthorType(val string) *Context { - self.Attributes[AuthorTypeKey] = val - return self -} - func (self *Context) SetTraceId(val string) *Context { self.Attributes[TraceIdKey] = val return self } -func (self *Context) SetSource(val string) *Context { - self.Attributes[Source] = val +func (self *Context) SetSourceType(val string) *Context { + self.Attributes[SourceType] = val return self } +func (self *Context) SetSourceAuth(val string) *Context { + self.Attributes[SourceAuth] = val + return self +} + +func (self *Context) SetSourceMethod(val string) *Context { + self.Attributes[SourceMethod] = val + return self +} + +func (self *Context) SetSourceLocal(val string) *Context { + self.Attributes[SourceLocal] = val + return self +} + +func (self *Context) SetSourceRemote(val string) *Context { + self.Attributes[SourceRemote] = val + return self +} + +func (self *Context) GetAuthor() *Author { + if self == nil { + return nil + } + return &Author{ + Type: self.Attributes[AuthorTypeKey], + Id: self.Attributes[AuthorIdKey], + Name: self.Attributes[AuthorNameKey], + } +} + +func (self *Context) GetSource() *Source { + if self == nil { + return nil + } + return &Source{ + Type: self.Attributes[SourceType], + Auth: self.Attributes[SourceAuth], + LocalAddr: self.Attributes[SourceLocal], + RemoteAddr: self.Attributes[SourceRemote], + Method: self.Attributes[SourceMethod], + } +} + +func (self *Context) PopulateMetadata(meta map[string]any) { + meta["author"] = self.GetAuthor() + meta["source"] = self.GetSource() + if traceId, found := self.Attributes[TraceIdKey]; found { + meta["trace_id"] = traceId + } +} + func (self *Context) ToProtoBuf() *cmd_pb.ChangeContext { if self == nil { return nil } + return &cmd_pb.ChangeContext{ Attributes: self.Attributes, RaftIndex: self.RaftIndex, diff --git a/controller/command/command.go b/controller/command/command.go index de24c60a0..732d1c901 100644 --- a/controller/command/command.go +++ b/controller/command/command.go @@ -78,6 +78,11 @@ func (self *LocalDispatcher) Dispatch(command Command) error { } }() + changeCtx := command.GetChangeContext() + if changeCtx == nil { + changeCtx = change.New().SetSourceType("unattributed").SetChangeAuthorType(change.AuthorTypeUnattributed) + } + ctx := changeCtx.NewMutateContext() if self.EncodeDecodeCommands { bytes, err := command.Encode() if err != nil { @@ -87,10 +92,10 @@ func (self *LocalDispatcher) Dispatch(command Command) error { if err != nil { return err } - return cmd.Apply(change.New().NewMutateContext()) + return cmd.Apply(ctx) } - return command.Apply(change.New().NewMutateContext()) + return command.Apply(ctx) } // Decoder instances know how to decode encoded commands diff --git a/controller/db/db.go b/controller/db/db.go index e30c6175f..96e47de20 100644 --- a/controller/db/db.go +++ b/controller/db/db.go @@ -21,7 +21,8 @@ import ( ) const ( - RootBucket = "ziti" + RootBucket = "ziti" + MetadataBucket = "metadata" ) func Open(path string) (boltz.Db, error) { diff --git a/controller/db/router_store.go b/controller/db/router_store.go index 69932022d..9e97984cb 100644 --- a/controller/db/router_store.go +++ b/controller/db/router_store.go @@ -32,11 +32,11 @@ const ( type Router struct { boltz.BaseExtEntity - Name string - Fingerprint *string - Cost uint16 - NoTraversal bool - Disabled bool + Name string `json:"name"` + Fingerprint *string `json:"fingerprint"` + Cost uint16 `json:"cost"` + NoTraversal bool `json:"no_traversal"` + Disabled bool `json:"disabled"` } func (entity *Router) GetEntityType() string { diff --git a/controller/db/service_store.go b/controller/db/service_store.go index cb7e93740..f2444998b 100644 --- a/controller/db/service_store.go +++ b/controller/db/service_store.go @@ -31,8 +31,8 @@ const ( type Service struct { boltz.BaseExtEntity - Name string - TerminatorStrategy string + Name string `json:"name"` + TerminatorStrategy string `json:"terminator_strategy"` } func (entity *Service) GetEntityType() string { diff --git a/controller/db/terminator_store.go b/controller/db/terminator_store.go index f81ae7218..b167f3fff 100644 --- a/controller/db/terminator_store.go +++ b/controller/db/terminator_store.go @@ -42,16 +42,16 @@ const ( type Terminator struct { boltz.BaseExtEntity - Service string - Router string - Binding string - Address string - InstanceId string - InstanceSecret []byte - Cost uint16 - Precedence string - PeerData xt.PeerData - HostId string + Service string `json:"service"` + Router string `json:"router"` + Binding string `json:"binding"` + Address string `json:"address"` + InstanceId string `json:"instance_id"` + InstanceSecret []byte `json:"instance_secret"` + Cost uint16 `json:"cost"` + Precedence string `json:"precedence"` + PeerData xt.PeerData `json:"peer_data"` + HostId string `json:"host_id"` } func (entity *Terminator) GetCost() uint16 { diff --git a/controller/handler_ctrl/base.go b/controller/handler_ctrl/base.go index b36741406..279f1477e 100644 --- a/controller/handler_ctrl/base.go +++ b/controller/handler_ctrl/base.go @@ -17,7 +17,6 @@ package handler_ctrl import ( - "fmt" "github.com/openziti/channel/v2" "github.com/openziti/fabric/controller/change" "github.com/openziti/fabric/controller/network" @@ -28,10 +27,13 @@ type baseHandler struct { network *network.Network } -func (self *baseHandler) newChangeContext(ch channel.Channel) *change.Context { +func (self *baseHandler) newChangeContext(ch channel.Channel, method string) *change.Context { return change.New(). SetChangeAuthorId(self.router.Id). SetChangeAuthorName(self.router.Name). - SetChangeAuthorType("router"). - SetSource(fmt.Sprintf("ctrl[%v]", ch.Underlay().GetRemoteAddr().String())) + SetChangeAuthorType(change.AuthorTypeRouter). + SetSourceType(change.SourceTypeControlChannel). + SetSourceMethod(method). + SetSourceLocal(ch.Underlay().GetLocalAddr().String()). + SetSourceRemote(ch.Underlay().GetRemoteAddr().String()) } diff --git a/controller/handler_ctrl/create_terminator.go b/controller/handler_ctrl/create_terminator.go index 6a9396428..e5b5fc9e8 100644 --- a/controller/handler_ctrl/create_terminator.go +++ b/controller/handler_ctrl/create_terminator.go @@ -72,7 +72,7 @@ func (self *createTerminatorHandler) handleCreateTerminator(msg *channel.Message Cost: uint16(request.Cost), } - if err := self.network.Terminators.Create(terminator, self.newChangeContext(ch)); err == nil { + if err := self.network.Terminators.Create(terminator, self.newChangeContext(ch, "fabric.create.terminator")); err == nil { pfxlog.Logger().Infof("created terminator [t/%s]", terminator.Id) handler_common.SendSuccess(msg, ch, terminator.Id) } else { diff --git a/controller/handler_ctrl/remove_terminator.go b/controller/handler_ctrl/remove_terminator.go index d9b411c18..dc8e95294 100644 --- a/controller/handler_ctrl/remove_terminator.go +++ b/controller/handler_ctrl/remove_terminator.go @@ -63,7 +63,7 @@ func (self *removeTerminatorHandler) handleRemoveTerminator(msg *channel.Message return } - if err := self.network.Terminators.Delete(request.TerminatorId, self.newChangeContext(ch)); err == nil { + if err := self.network.Terminators.Delete(request.TerminatorId, self.newChangeContext(ch, "fabric.remove.terminator")); err == nil { log. WithField("routerId", ch.Id()). WithField("serviceId", terminator.Service). diff --git a/controller/handler_ctrl/remove_terminators.go b/controller/handler_ctrl/remove_terminators.go index 2a537f7c1..5aa1d7cf2 100644 --- a/controller/handler_ctrl/remove_terminators.go +++ b/controller/handler_ctrl/remove_terminators.go @@ -57,7 +57,7 @@ func (self *removeTerminatorsHandler) HandleReceive(msg *channel.Message, ch cha func (self *removeTerminatorsHandler) handleRemoveTerminators(msg *channel.Message, ch channel.Channel, request *ctrl_pb.RemoveTerminatorsRequest) { log := pfxlog.ContextLogger(ch.Label()) - if err := self.network.Terminators.DeleteBatch(request.TerminatorIds, self.newChangeContext(ch)); err == nil { + if err := self.network.Terminators.DeleteBatch(request.TerminatorIds, self.newChangeContext(ch, "fabric.remove.terminators.batch")); err == nil { log. WithField("routerId", ch.Id()). WithField("terminatorIds", request.TerminatorIds). diff --git a/controller/handler_ctrl/update_terminator.go b/controller/handler_ctrl/update_terminator.go index 2fbe30ce6..72439516e 100644 --- a/controller/handler_ctrl/update_terminator.go +++ b/controller/handler_ctrl/update_terminator.go @@ -97,7 +97,7 @@ func (self *updateTerminatorHandler) handleUpdateTerminator(msg *channel.Message checker[db.FieldTerminatorPrecedence] = struct{}{} } - if err := self.network.Terminators.Update(terminator, checker, self.newChangeContext(ch)); err != nil { + if err := self.network.Terminators.Update(terminator, checker, self.newChangeContext(ch, "fabric.update.terminator")); err != nil { handler_common.SendFailure(msg, ch, err.Error()) return } diff --git a/controller/network/routesender.go b/controller/network/routesender.go index 8fe87f8ef..e46ec1f6a 100644 --- a/controller/network/routesender.go +++ b/controller/network/routesender.go @@ -165,8 +165,15 @@ func (self *routeSender) handleRouteSend(attempt uint32, path *Path, strategy xt changeCtx := change.New(). SetChangeAuthorId(status.Router.Id). SetChangeAuthorName(status.Router.Name). - SetChangeAuthorType("router"). - SetSource("router reported invalid terminator") + SetChangeAuthorType(change.AuthorTypeRouter). + SetSourceType(change.SourceTypeControlChannel). + SetSourceMethod("route.response") + if ch := status.Router.Control; ch != nil { + changeCtx. + SetSourceLocal(ch.Underlay().GetLocalAddr().String()). + SetSourceRemote(ch.Underlay().GetRemoteAddr().String()) + } + if err := self.terminators.Delete(terminator.GetId(), changeCtx); err != nil { logger.WithError(fmt.Errorf("unable to delete invalid terminator: %v", err)) } diff --git a/controller/network/terminator.go b/controller/network/terminator.go index b0e8d5200..7955f96ef 100644 --- a/controller/network/terminator.go +++ b/controller/network/terminator.go @@ -196,7 +196,7 @@ func (self *TerminatorManager) handlePrecedenceChange(terminatorId string, prece db.FieldTerminatorPrecedence: struct{}{}, } - if err = self.Update(terminator, checker, change.New().SetChangeAuthorId("xt").SetChangeAuthorType("controller")); err != nil { + if err = self.Update(terminator, checker, change.New().SetSourceType(change.SourceTypeXt).SetChangeAuthorType(change.AuthorTypeController)); err != nil { pfxlog.Logger().Errorf("unable to update precedence for terminator %v to %v (%v)", terminatorId, precedence, err) } } diff --git a/controller/raft/fsm.go b/controller/raft/fsm.go index 393d70263..a05d2d11f 100644 --- a/controller/raft/fsm.go +++ b/controller/raft/fsm.go @@ -35,8 +35,7 @@ import ( ) const ( - bucketName = "raft" - fieldIndex = "index" + fieldIndex = "raftIndex" ) func NewFsm(dataDir string, decoders command.Decoders, indexTracker IndexTracker, eventDispatcher event.Dispatcher) *BoltDbFsm { @@ -85,7 +84,7 @@ func (self *BoltDbFsm) GetDb() boltz.Db { func (self *BoltDbFsm) loadCurrentIndex() (uint64, error) { var result uint64 err := self.db.View(func(tx *bbolt.Tx) error { - if raftBucket := boltz.Path(tx, db.RootBucket, bucketName); raftBucket != nil { + if raftBucket := boltz.Path(tx, db.RootBucket, db.MetadataBucket); raftBucket != nil { if val := raftBucket.GetInt64(fieldIndex); val != nil { result = uint64(*val) } @@ -96,7 +95,7 @@ func (self *BoltDbFsm) loadCurrentIndex() (uint64, error) { } func (self *BoltDbFsm) updateIndexInTx(tx *bbolt.Tx, index uint64) error { - raftBucket := boltz.GetOrCreatePath(tx, db.RootBucket, bucketName) + raftBucket := boltz.GetOrCreatePath(tx, db.RootBucket, db.MetadataBucket) raftBucket.SetInt64(fieldIndex, int64(index), nil) return raftBucket.GetError() } @@ -156,7 +155,7 @@ func (self *BoltDbFsm) Apply(log *raft.Log) interface{} { logger.Infof("apply log with type %T", cmd) changeCtx := cmd.GetChangeContext() if changeCtx == nil { - changeCtx = change.New().SetSource("untracked") + changeCtx = change.New().SetSourceType("unattributed").SetChangeAuthorType(change.AuthorTypeUnattributed) } changeCtx.RaftIndex = log.Index diff --git a/event/dispatcher.go b/event/dispatcher.go index 3f61834b2..81cb869cc 100644 --- a/event/dispatcher.go +++ b/event/dispatcher.go @@ -17,6 +17,7 @@ package event import ( + "github.com/openziti/storage/boltz" "io" "regexp" ) @@ -109,7 +110,13 @@ type Dispatcher interface { AddClusterEventHandler(handler ClusterEventHandler) RemoveClusterEventHandler(handler ClusterEventHandler) + AddEntityChangeEventHandler(handler EntityChangeEventHandler) + RemoveEntityChangeEventHandler(handler EntityChangeEventHandler) + AddEntityChangeSource(store boltz.Store) + AddGlobalEntityChangeMetadata(k string, v any) + CircuitEventHandler + EntityChangeEventHandler LinkEventHandler MetricsEventHandler MetricsMessageHandler diff --git a/event/dispatcher_mock.go b/event/dispatcher_mock.go index d05c9002d..92ceb8253 100644 --- a/event/dispatcher_mock.go +++ b/event/dispatcher_mock.go @@ -18,6 +18,7 @@ package event import ( "github.com/openziti/metrics/metrics_pb" + "github.com/openziti/storage/boltz" "regexp" ) @@ -25,6 +26,16 @@ var _ Dispatcher = DispatcherMock{} type DispatcherMock struct{} +func (d DispatcherMock) AddEntityChangeSource(store boltz.Store) {} + +func (d DispatcherMock) AddGlobalEntityChangeMetadata(k string, v any) {} + +func (d DispatcherMock) AddEntityChangeEventHandler(handler EntityChangeEventHandler) {} + +func (d DispatcherMock) RemoveEntityChangeEventHandler(handler EntityChangeEventHandler) {} + +func (d DispatcherMock) AcceptEntityChangeEvent(event *EntityChangeEvent) {} + func (d DispatcherMock) GetFormatterFactory(formatterType string) FormatterFactory { return nil } diff --git a/event/entity_change.go b/event/entity_change.go new file mode 100644 index 000000000..c043e3027 --- /dev/null +++ b/event/entity_change.go @@ -0,0 +1,55 @@ +/* + 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 event + +import ( + "time" +) + +type EntityChangeEventType string + +const ( + EntityChangeEventsNs = "entityChange" + + EntityChangeTypeEntityCreated EntityChangeEventType = "created" + EntityChangeTypeEntityUpdated EntityChangeEventType = "updated" + EntityChangeTypeEntityDeleted EntityChangeEventType = "deleted" + EntityChangeTypeCommitted EntityChangeEventType = "committed" +) + +type EntityChangeEvent struct { + Namespace string `json:"namespace"` + EventId string `json:"event_id"` + EventType EntityChangeEventType `json:"event_type"` + Timestamp time.Time `json:"timestamp"` + Metadata map[string]any `json:"metadata,omitempty"` + EntityType string `json:"entity_type,omitempty"` + IsParentEvent *bool `json:"is_parent_event,omitempty"` + InitialState any `json:"initial_state,omitempty"` + FinalState any `json:"final_state,omitempty"` + PropagateIndicator bool `json:"-"` + IsRecoveryEvent bool `json:"-"` +} + +type EntityChangeEventHandler interface { + AcceptEntityChangeEvent(event *EntityChangeEvent) +} + +type EntityChangeEventHandlerWrapper interface { + EntityChangeEventHandler + IsWrapping(value EntityChangeEventHandler) bool +} diff --git a/events/dispatcher.go b/events/dispatcher.go index 9f4a21ffa..95dd4e3f6 100644 --- a/events/dispatcher.go +++ b/events/dispatcher.go @@ -45,9 +45,15 @@ func NewDispatcher(closeNotify <-chan struct{}) *Dispatcher { result := &Dispatcher{ closeNotify: closeNotify, eventC: make(chan event.Event, 25), + entityChangeEventsDispatcher: entityChangeEventDispatcher{ + notifyCh: make(chan struct{}, 1), + globalMetadata: map[string]any{}, + }, } + result.entityChangeEventsDispatcher.dispatcher = result result.RegisterEventTypeFunctions(event.CircuitEventsNs, result.registerCircuitEventHandler, result.unregisterCircuitEventHandler) + result.RegisterEventTypeFunctions(event.EntityChangeEventsNs, result.registerEntityChangeEventHandler, result.unregisterEntityChangeEventHandler) result.RegisterEventTypeFunctions(event.LinkEventsNs, result.registerLinkEventHandler, result.unregisterLinkEventHandler) result.RegisterEventTypeFunctions(event.MetricsEventsNs, result.registerMetricsEventHandler, result.unregisterMetricsEventHandler) result.RegisterEventTypeFunctions(event.RouterEventsNs, result.registerRouterEventHandler, result.unregisterRouterEventHandler) @@ -72,16 +78,17 @@ func NewDispatcher(closeNotify <-chan struct{}) *Dispatcher { var _ event.Dispatcher = (*Dispatcher)(nil) type Dispatcher struct { - circuitEventHandlers concurrenz.CopyOnWriteSlice[event.CircuitEventHandler] - linkEventHandlers concurrenz.CopyOnWriteSlice[event.LinkEventHandler] - metricsEventHandlers concurrenz.CopyOnWriteSlice[event.MetricsEventHandler] - metricsMsgEventHandlers concurrenz.CopyOnWriteSlice[event.MetricsMessageHandler] - routerEventHandlers concurrenz.CopyOnWriteSlice[event.RouterEventHandler] - serviceEventHandlers concurrenz.CopyOnWriteSlice[event.ServiceEventHandler] - terminatorEventHandlers concurrenz.CopyOnWriteSlice[event.TerminatorEventHandler] - usageEventHandlers concurrenz.CopyOnWriteSlice[event.UsageEventHandler] - usageEventV3Handlers concurrenz.CopyOnWriteSlice[event.UsageEventV3Handler] - clusterEventHandlers concurrenz.CopyOnWriteSlice[event.ClusterEventHandler] + circuitEventHandlers concurrenz.CopyOnWriteSlice[event.CircuitEventHandler] + entityChangeEventHandlers concurrenz.CopyOnWriteSlice[event.EntityChangeEventHandler] + linkEventHandlers concurrenz.CopyOnWriteSlice[event.LinkEventHandler] + metricsEventHandlers concurrenz.CopyOnWriteSlice[event.MetricsEventHandler] + metricsMsgEventHandlers concurrenz.CopyOnWriteSlice[event.MetricsMessageHandler] + routerEventHandlers concurrenz.CopyOnWriteSlice[event.RouterEventHandler] + serviceEventHandlers concurrenz.CopyOnWriteSlice[event.ServiceEventHandler] + terminatorEventHandlers concurrenz.CopyOnWriteSlice[event.TerminatorEventHandler] + usageEventHandlers concurrenz.CopyOnWriteSlice[event.UsageEventHandler] + usageEventV3Handlers concurrenz.CopyOnWriteSlice[event.UsageEventV3Handler] + clusterEventHandlers concurrenz.CopyOnWriteSlice[event.ClusterEventHandler] metricsMappers concurrenz.CopyOnWriteSlice[event.MetricsMapper] @@ -89,8 +96,10 @@ type Dispatcher struct { eventHandlerFactories concurrenz.CopyOnWriteMap[string, event.HandlerFactory] formatterFactories concurrenz.CopyOnWriteMap[string, event.FormatterFactory] - closeNotify <-chan struct{} - eventC chan event.Event + entityChangeEventsDispatcher entityChangeEventDispatcher + entityTypes []string + closeNotify <-chan struct{} + eventC chan event.Event } func (self *Dispatcher) InitializeNetworkEvents(n *network.Network) { @@ -99,6 +108,7 @@ func (self *Dispatcher) InitializeNetworkEvents(n *network.Network) { self.initServiceEvents(n) self.initTerminatorEvents(n) self.initUsageEvents() + self.initEntityChangeEvents(n) self.AddMetricsMapper(ctrlChannelMetricsMapper{}.mapMetrics) self.AddMetricsMapper((&linkMetricsMapper{network: n}).mapMetrics) @@ -188,6 +198,7 @@ func (self *Dispatcher) WireEventHandlers(eventHandlerConfigs []*EventHandlerCon return err } } + self.entityChangeEventsDispatcher.flushCommittedTxEvents(true) return nil } diff --git a/events/dispatcher_entity_change.go b/events/dispatcher_entity_change.go new file mode 100644 index 000000000..ba2105abe --- /dev/null +++ b/events/dispatcher_entity_change.go @@ -0,0 +1,367 @@ +/* + 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 events + +import ( + "context" + "encoding/binary" + "github.com/michaelquigley/pfxlog" + "github.com/openziti/fabric/controller/change" + "github.com/openziti/fabric/controller/db" + "github.com/openziti/fabric/controller/network" + "github.com/openziti/fabric/event" + "github.com/openziti/foundation/v2/genext" + "github.com/openziti/storage/boltz" + "github.com/pkg/errors" + "go.etcd.io/bbolt" + "reflect" + "strings" + "time" +) + +const ( + entityChangeEventsBucket = "entityChangeEvents" +) + +func (self *Dispatcher) AddEntityChangeEventHandler(handler event.EntityChangeEventHandler) { + self.entityChangeEventHandlers.Append(handler) +} + +func (self *Dispatcher) RemoveEntityChangeEventHandler(handler event.EntityChangeEventHandler) { + self.entityChangeEventHandlers.DeleteIf(func(val event.EntityChangeEventHandler) bool { + if val == handler { + return true + } + if w, ok := val.(event.EntityChangeEventHandlerWrapper); ok { + return w.IsWrapping(handler) + } + return false + }) +} + +func (self *Dispatcher) AcceptEntityChangeEvent(event *event.EntityChangeEvent) { + // don't do these in a separate goroutine to minimize the chance of losing events + // If we need to, the handler can spin up a separate goroutine + for _, handler := range self.entityChangeEventHandlers.Value() { + handler.AcceptEntityChangeEvent(event) + } +} + +func (self *Dispatcher) registerEntityChangeEventHandler(val interface{}, options map[string]interface{}) error { + handler, ok := val.(event.EntityChangeEventHandler) + + if !ok { + return errors.Errorf("type %v doesn't implement github.com/openziti/fabric/event/EntityChangeEventHandler interface.", reflect.TypeOf(val)) + } + + propagateAlways := false + if val, found := options["propagateAlways"]; found { + if b, ok := val.(bool); ok { + propagateAlways = b + } else if s, ok := val.(string); ok { + propagateAlways = strings.EqualFold(s, "true") + } else { + return errors.New("invalid value for entityChange.propagateAlways, must be boolean or string") + } + } + + includeParentEvents := false + if val, found := options["includeParentEvents"]; found { + if b, ok := val.(bool); ok { + includeParentEvents = b + } else if s, ok := val.(string); ok { + includeParentEvents = strings.EqualFold(s, "true") + } else { + return errors.New("invalid value for entityChange.includeParentEvents, must be boolean or string") + } + } + + filter := &entityChangeEventFilter{ + EntityChangeEventHandler: handler, + propagateAlways: propagateAlways, + includeParentEvents: includeParentEvents, + } + + if val, found := options["include"]; found { + includes := map[string]struct{}{} + if list, ok := val.([]interface{}); ok { + for _, val := range list { + if entityType, ok := val.(string); ok { + includes[entityType] = struct{}{} + } else { + return errors.Errorf("invalid value type [%T] for entityChange include list, must be string list", val) + } + } + } else { + return errors.Errorf("invalid value type [%T] for entityChange include list, must be string list", val) + } + + if len(includes) == 0 { + return errors.Errorf("no values provided in include list for entityChange events, either drop includes stanza or provide at least one entity type to include") + } + + for entityType := range includes { + if !genext.Contains(self.entityTypes, entityType) { + return errors.Errorf("invalid entity type [%v] in entityChange events include list, valid values include: %v", entityType, self.entityTypes) + } + } + + filter.entityTypes = includes + } + + self.AddEntityChangeEventHandler(filter) + + return nil +} + +func (self *Dispatcher) unregisterEntityChangeEventHandler(val interface{}) { + if handler, ok := val.(event.EntityChangeEventHandler); ok { + self.RemoveEntityChangeEventHandler(handler) + } +} + +func (self *Dispatcher) initEntityChangeEvents(n *network.Network) { + self.entityChangeEventsDispatcher.network = n + for _, store := range n.GetStores().GetStoreList() { + self.AddEntityChangeSource(store) + } + self.AddGlobalEntityChangeMetadata("version", n.VersionProvider.Version()) + go self.entityChangeEventsDispatcher.flushLoop() +} + +func (self *Dispatcher) AddEntityChangeSource(store boltz.Store) { + store.AddUntypedEntityConstraint(&self.entityChangeEventsDispatcher) + self.entityTypes = append(self.entityTypes, store.GetEntityType()) +} + +func (self *Dispatcher) AddGlobalEntityChangeMetadata(k string, v any) { + self.entityChangeEventsDispatcher.globalMetadata[k] = v +} + +func txIdToBytes(txId uint64) []byte { + return binary.LittleEndian.AppendUint64(nil, txId) +} + +func bytesToTxId(b []byte) uint64 { + return binary.LittleEndian.Uint64(b) +} + +type entityChangeEventDispatcher struct { + network *network.Network + dispatcher *Dispatcher + notifyCh chan struct{} + globalMetadata map[string]any +} + +func (self *entityChangeEventDispatcher) logTxEvent(state boltz.UntypedEntityChangeState) error { + tx := state.GetCtx().Tx() + rowId := txIdToBytes(uint64(tx.ID())) + eventsBucket := boltz.GetOrCreatePath(tx, db.RootBucket, db.MetadataBucket, entityChangeEventsBucket) + txBucket, err := eventsBucket.CreateBucketIfNotExists(rowId) + if err != nil { + return err + } + return txBucket.Put([]byte(state.GetEventId()), []byte(state.GetStore().GetEntityType())) +} + +func (self *entityChangeEventDispatcher) processPreviousTxEvents(tx *bbolt.Tx, emit bool) { + currentTxId := uint64(tx.ID()) + + eventsBucket := boltz.GetOrCreatePath(tx, db.RootBucket, db.MetadataBucket, entityChangeEventsBucket) + txCursor := eventsBucket.OpenCursor(tx, true) + + var toDelete []uint64 + for txCursor.IsValid() { + rowId := txCursor.Current() + txId := bytesToTxId(rowId) + log := pfxlog.Logger().WithField("txId", txId) + if txId != currentTxId { + log.Debug("cleaning up entity change events for tx") + + if emit { + txBucket := eventsBucket.GetBucketByKey(rowId) + eventIdCursor := txBucket.Cursor() + for k, v := eventIdCursor.First(); k != nil; k, v = eventIdCursor.Next() { + eventId := string(k) + entityType := string(v) + log.WithField("eventId", eventId).WithField("entityType", entityType).Debug("emitting event for tx") + self.emitRecoveryEvent(eventId, entityType) + eventIdCursor.Next() + } + } + toDelete = append(toDelete, txId) + } + txCursor.Next() + } + + if len(toDelete) > 0 { + for _, txId := range toDelete { + if err := eventsBucket.DeleteBucket(txIdToBytes(txId)); err != nil { + pfxlog.Logger().WithError(err).WithField("txId", txId).Error("unable to delete event bucket for tx") + } + } + } +} + +func (self *entityChangeEventDispatcher) ProcessPreCommit(state boltz.UntypedEntityChangeState) error { + self.processPreviousTxEvents(state.GetCtx().Tx(), false) + + var changeType event.EntityChangeEventType + if state.GetChangeType() == boltz.EntityCreated { + changeType = event.EntityChangeTypeEntityCreated + } else if state.GetChangeType() == boltz.EntityUpdated { + changeType = event.EntityChangeTypeEntityUpdated + } else if state.GetChangeType() == boltz.EntityDeleted { + changeType = event.EntityChangeTypeEntityDeleted + } + + isParentEvent := state.IsParentEvent() + evt := &event.EntityChangeEvent{ + Namespace: event.EntityChangeEventsNs, + EventId: state.GetEventId(), + EventType: changeType, + EntityType: state.GetStore().GetEntityType(), + IsParentEvent: &isParentEvent, + Timestamp: time.Now(), + Metadata: map[string]any{}, + InitialState: state.GetInitialState(), + FinalState: state.GetFinalState(), + PropagateIndicator: self.network.Dispatcher.IsLeaderOrLeaderless(), + } + + changeCtx := change.FromContext(state.GetCtx().Context()) + if changeCtx == nil { + changeCtx = change.New() + state.GetCtx().UpdateContext(func(ctx context.Context) context.Context { + return changeCtx.AddToContext(ctx) + }) + } + + changeCtx.PopulateMetadata(evt.Metadata) + + for k, v := range self.globalMetadata { + evt.Metadata[k] = v + } + + self.dispatcher.AcceptEntityChangeEvent(evt) + return self.logTxEvent(state) +} + +func (self *entityChangeEventDispatcher) emitRecoveryEvent(eventId string, entityType string) { + evt := &event.EntityChangeEvent{ + Namespace: event.EntityChangeEventsNs, + EventId: eventId, + EntityType: entityType, + EventType: event.EntityChangeTypeCommitted, + Timestamp: time.Now(), + IsRecoveryEvent: true, + } + self.dispatcher.AcceptEntityChangeEvent(evt) +} + +func (self *entityChangeEventDispatcher) ProcessPostCommit(state boltz.UntypedEntityChangeState) { + isParentEvent := state.IsParentEvent() + evt := &event.EntityChangeEvent{ + Namespace: event.EntityChangeEventsNs, + EventId: state.GetEventId(), + EventType: event.EntityChangeTypeCommitted, + EntityType: state.GetStore().GetEntityType(), + Timestamp: time.Now(), + IsParentEvent: &isParentEvent, + PropagateIndicator: self.network.Dispatcher.IsLeaderOrLeaderless(), + } + self.dispatcher.AcceptEntityChangeEvent(evt) + self.notifyFlush() +} + +func (self *entityChangeEventDispatcher) notifyFlush() { + select { + case self.notifyCh <- struct{}{}: + default: + } +} + +func (self *entityChangeEventDispatcher) flushLoop() { + for { + // wait to be notified of an event + <-self.notifyCh + + // wait until we've not gotten an event for 5 seconds before cleaning up + flushed := false + for !flushed { + select { + case <-self.notifyCh: + case <-time.After(5 * time.Second): + pfxlog.Logger().Debug("cleaning up entity change events") + self.flushCommittedTxEvents(false) + flushed = true + } + } + } +} + +func (self *entityChangeEventDispatcher) flushCommittedTxEvents(emit bool) { + err := self.network.GetDb().Update(nil, func(ctx boltz.MutateContext) error { + self.processPreviousTxEvents(ctx.Tx(), emit) + return nil + }) + if err != nil { + pfxlog.Logger().WithError(err).Error("error while flushing committed tx entity change events") + } +} + +type entityChangeEventFilter struct { + event.EntityChangeEventHandler + propagateAlways bool + includeParentEvents bool + entityTypes map[string]struct{} +} + +func (self *entityChangeEventFilter) IsWrapping(value event.EntityChangeEventHandler) bool { + if self.EntityChangeEventHandler == value { + return true + } + if w, ok := self.EntityChangeEventHandler.(event.EntityChangeEventHandlerWrapper); ok { + return w.IsWrapping(value) + } + return false +} + +func (self *entityChangeEventFilter) AcceptEntityChangeEvent(evt *event.EntityChangeEvent) { + if !evt.IsRecoveryEvent { + if !self.propagateAlways && !evt.PropagateIndicator { + return + } + + if *evt.IsParentEvent && !self.includeParentEvents { + return + } + } + + if self.entityTypes != nil { + if _, found := self.entityTypes[evt.EntityType]; !found { + return + } + } + + if evt.EventType == event.EntityChangeTypeCommitted { + evt.IsParentEvent = nil + evt.EntityType = "" + } + + self.EntityChangeEventHandler.AcceptEntityChangeEvent(evt) +} diff --git a/events/formatter.go b/events/formatter.go index 1ef608f37..71ef48aeb 100644 --- a/events/formatter.go +++ b/events/formatter.go @@ -185,6 +185,16 @@ func (event *JsonClusterEvent) Format() ([]byte, error) { return MarshalJson(event) } +type JsonEntityChangeEvent event.EntityChangeEvent + +func (event *JsonEntityChangeEvent) GetEventType() string { + return "entity.change" +} + +func (event *JsonEntityChangeEvent) Format() ([]byte, error) { + return MarshalJson(event) +} + func NewJsonFormatter(queueDepth int, sink event.FormattedEventSink) *JsonFormatter { result := &JsonFormatter{ BaseFormatter: BaseFormatter{ @@ -237,6 +247,10 @@ func (formatter *JsonFormatter) AcceptClusterEvent(evt *event.ClusterEvent) { formatter.AcceptLoggingEvent((*JsonClusterEvent)(evt)) } +func (formatter *JsonFormatter) AcceptEntityChangeEvent(evt *event.EntityChangeEvent) { + formatter.AcceptLoggingEvent((*JsonEntityChangeEvent)(evt)) +} + var histogramBuckets = map[string]string{"p50": "0.50", "p75": "0.75", "p95": "0.95", "p99": "0.99", "p999": "0.999", "p9999": "0.9999"} type PrometheusMetricsEvent event.MetricsEvent