mirror of
https://github.com/openziti/ziti.git
synced 2026-10-04 04:21:39 +00:00
248 lines
7.0 KiB
Go
248 lines
7.0 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 common
|
|
|
|
import (
|
|
"fmt"
|
|
"github.com/openziti/ziti/common/pb/edge_ctrl_pb"
|
|
"sync"
|
|
"sync/atomic"
|
|
)
|
|
|
|
type OnStoreSuccess func(index uint64, event *edge_ctrl_pb.DataState_ChangeSet)
|
|
|
|
type EventCache interface {
|
|
// Store allows storage of an event and execution of an onSuccess callback while the event cache remains locked.
|
|
// onSuccess may be nil. This function is blocking.
|
|
Store(event *edge_ctrl_pb.DataState_ChangeSet, onSuccess OnStoreSuccess) error
|
|
|
|
// CurrentIndex returns the latest event index applied. This function is blocking.
|
|
CurrentIndex() (uint64, bool)
|
|
|
|
// ReplayFrom returns an array of events from startIndex and true if the replay may be facilitated.
|
|
// An empty slice and true is returned in cases where the requested startIndex is greater than the current index.
|
|
// An empty slice and false is returned in cases where the replay cannot be facilitated.
|
|
// This function is blocking.
|
|
ReplayFrom(startIndex uint64) ([]*edge_ctrl_pb.DataState_ChangeSet, bool)
|
|
|
|
// WhileLocked allows the execution of arbitrary functionality while the event cache is locked. This function
|
|
// is blocking.
|
|
WhileLocked(func(uint64, bool))
|
|
|
|
// SetCurrentIndex sets the current index to the supplied value. All event log history may be lost.
|
|
SetCurrentIndex(uint64)
|
|
}
|
|
|
|
// ForgetfulEventCache does not store events or support replaying. It tracks
|
|
// the event index and that is it. It is a stand in for LoggingEventCache
|
|
// when replaying events is not expected (i.e. in routers)
|
|
type ForgetfulEventCache struct {
|
|
lock sync.Mutex
|
|
index uint64
|
|
}
|
|
|
|
func NewForgetfulEventCache() *ForgetfulEventCache {
|
|
return &ForgetfulEventCache{}
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) SetCurrentIndex(index uint64) {
|
|
cache.index = index
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) WhileLocked(callback func(uint64, bool)) {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
callback(cache.currentIndex())
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) Store(event *edge_ctrl_pb.DataState_ChangeSet, onSuccess OnStoreSuccess) error {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
// Synthetic events are not backed by any kind of data store that provides and index. They are not stored and
|
|
// trigger the on success callback immediately.
|
|
if event.IsSynthetic {
|
|
onSuccess(event.Index, event)
|
|
return nil
|
|
}
|
|
|
|
if cache.index > 0 && cache.index >= event.Index {
|
|
return fmt.Errorf("out of order event detected, currentIndex: %d, receivedIndex: %d, type :%T", cache.index, event.Index, cache)
|
|
}
|
|
|
|
cache.index = event.Index
|
|
|
|
if onSuccess != nil {
|
|
onSuccess(cache.index, event)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) ReplayFrom(_ uint64) ([]*edge_ctrl_pb.DataState_ChangeSet, bool) {
|
|
return nil, false
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) CurrentIndex() (uint64, bool) {
|
|
return cache.currentIndex()
|
|
}
|
|
|
|
func (cache *ForgetfulEventCache) currentIndex() (uint64, bool) {
|
|
return atomic.LoadUint64(&cache.index), true
|
|
}
|
|
|
|
// LoggingEventCache stores events in order to support replaying (i.e. in controllers).
|
|
type LoggingEventCache struct {
|
|
lock sync.Mutex
|
|
HeadLogIndex uint64 `json:"-"`
|
|
LogSize uint64 `json:"-"`
|
|
Log []uint64 `json:"-"`
|
|
Events map[uint64]*edge_ctrl_pb.DataState_ChangeSet `json:"-"`
|
|
}
|
|
|
|
func NewLoggingEventCache(logSize uint64) *LoggingEventCache {
|
|
return &LoggingEventCache{
|
|
HeadLogIndex: 0,
|
|
LogSize: logSize,
|
|
Log: make([]uint64, logSize),
|
|
Events: map[uint64]*edge_ctrl_pb.DataState_ChangeSet{},
|
|
}
|
|
}
|
|
|
|
func (cache *LoggingEventCache) SetCurrentIndex(index uint64) {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
cache.HeadLogIndex = 0
|
|
cache.Log = make([]uint64, cache.LogSize)
|
|
cache.Log[0] = index
|
|
cache.Events = map[uint64]*edge_ctrl_pb.DataState_ChangeSet{}
|
|
}
|
|
|
|
func (cache *LoggingEventCache) WhileLocked(callback func(uint64, bool)) {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
callback(cache.currentIndex())
|
|
}
|
|
|
|
func (cache *LoggingEventCache) Store(event *edge_ctrl_pb.DataState_ChangeSet, onSuccess OnStoreSuccess) error {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
// Synthetic events are not backed by any kind of data store that provides and index. They are not stored and
|
|
// trigger the on success callback immediately.
|
|
if event.IsSynthetic {
|
|
onSuccess(event.Index, event)
|
|
return nil
|
|
}
|
|
|
|
currentIndex, ok := cache.currentIndex()
|
|
|
|
if ok && currentIndex >= event.Index {
|
|
return fmt.Errorf("out of order event detected, currentIndex: %d, receivedIndex: %d, type :%T", currentIndex, event.Index, cache)
|
|
}
|
|
|
|
targetLogIndex := uint64(0)
|
|
targetLogIndex = (cache.HeadLogIndex + 1) % cache.LogSize
|
|
|
|
// delete old value if we have looped
|
|
prevKey := cache.Log[targetLogIndex]
|
|
|
|
if prevKey != 0 {
|
|
delete(cache.Events, prevKey)
|
|
}
|
|
|
|
// add new values
|
|
cache.Log[targetLogIndex] = event.Index
|
|
cache.Events[event.Index] = event
|
|
|
|
//update head
|
|
cache.HeadLogIndex = targetLogIndex
|
|
|
|
onSuccess(event.Index, event)
|
|
return nil
|
|
}
|
|
|
|
func (cache *LoggingEventCache) CurrentIndex() (uint64, bool) {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
return cache.currentIndex()
|
|
}
|
|
|
|
func (cache *LoggingEventCache) currentIndex() (uint64, bool) {
|
|
if len(cache.Log) == 0 {
|
|
return 0, false
|
|
}
|
|
|
|
return cache.Log[cache.HeadLogIndex], true
|
|
}
|
|
|
|
func (cache *LoggingEventCache) ReplayFrom(startIndex uint64) ([]*edge_ctrl_pb.DataState_ChangeSet, bool) {
|
|
cache.lock.Lock()
|
|
defer cache.lock.Unlock()
|
|
|
|
_, eventFound := cache.Events[startIndex]
|
|
|
|
if !eventFound {
|
|
// if we're asked to replay an index we haven't reached yet, return an empty list
|
|
headIndex := cache.Log[cache.HeadLogIndex]
|
|
if headIndex < startIndex {
|
|
return nil, true
|
|
}
|
|
|
|
return nil, false
|
|
}
|
|
|
|
var startLogIndex *uint64
|
|
|
|
for logIndex, eventIndex := range cache.Log {
|
|
if eventIndex == startIndex {
|
|
tmp := uint64(logIndex)
|
|
startLogIndex = &tmp
|
|
break
|
|
}
|
|
}
|
|
|
|
if startLogIndex == nil {
|
|
return nil, false
|
|
}
|
|
|
|
// replay, no loop required
|
|
if *startLogIndex <= cache.HeadLogIndex {
|
|
var result []*edge_ctrl_pb.DataState_ChangeSet
|
|
for _, key := range cache.Log[*startLogIndex : cache.HeadLogIndex+1] {
|
|
result = append(result, cache.Events[key])
|
|
}
|
|
return result, true
|
|
}
|
|
|
|
// looping replay
|
|
var result []*edge_ctrl_pb.DataState_ChangeSet
|
|
for _, key := range cache.Log[*startLogIndex:] {
|
|
result = append(result, cache.Events[key])
|
|
}
|
|
|
|
for _, key := range cache.Log[:cache.HeadLogIndex+1] {
|
|
result = append(result, cache.Events[key])
|
|
}
|
|
|
|
return result, true
|
|
}
|