Files
fission-src/pkg/utils/namespace_manager.go
T

242 lines
6.2 KiB
Go

package utils
import (
"sort"
"sync"
"time"
)
type NamespaceSubscriber interface {
Name() string
}
type NamespaceManager interface {
Snapshot() []string
SnapshotRecords() []NamespaceRecord
Get(name string) (NamespaceRecord, bool)
Bootstrap(namespaces []string, source NamespaceSource, observedAt time.Time) []NamespaceRecord
Subscribe(subscriber NamespaceSubscriber)
SnapshotSubscribers() []string
Upsert(event NamespaceEvent) NamespaceRecord
MarkPartState(namespace string, part string, state NamespacePartState) (NamespaceRecord, bool)
MarkPartRegistering(namespace string, part string) (NamespaceRecord, bool)
MarkPartActive(namespace string, part string) (NamespaceRecord, bool)
MarkPartFailed(namespace string, part string, err error) (NamespaceRecord, bool)
Remove(name string) bool
}
type inMemoryNamespaceManager struct {
mu sync.RWMutex
records map[string]NamespaceRecord
subs map[string]NamespaceSubscriber
}
func NewNamespaceManager() NamespaceManager {
return &inMemoryNamespaceManager{
records: make(map[string]NamespaceRecord),
subs: make(map[string]NamespaceSubscriber),
}
}
func NewBootstrappedNamespaceManager(resolver *NamespaceResolver, source NamespaceSource, observedAt time.Time) NamespaceManager {
manager := NewNamespaceManager()
if resolver == nil {
return manager
}
manager.Bootstrap(resolver.Snapshot(), source, observedAt)
return manager
}
func (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) {
m.mu.Lock()
defer m.mu.Unlock()
m.subs[subscriber.Name()] = subscriber
}
func (m *inMemoryNamespaceManager) Bootstrap(namespaces []string, source NamespaceSource, observedAt time.Time) []NamespaceRecord {
records := make([]NamespaceRecord, 0, len(namespaces))
for _, namespace := range namespaces {
records = append(records, m.Upsert(NamespaceEvent{
Type: NamespaceEventAdd,
Name: namespace,
Source: source,
ObservedAt: observedAt,
}))
}
sort.Slice(records, func(i, j int) bool {
return records[i].Name < records[j].Name
})
return records
}
func (m *inMemoryNamespaceManager) SnapshotSubscribers() []string {
m.mu.RLock()
defer m.mu.RUnlock()
names := make([]string, 0, len(m.subs))
for name := range m.subs {
names = append(names, name)
}
sort.Strings(names)
return names
}
func (m *inMemoryNamespaceManager) Snapshot() []string {
m.mu.RLock()
defer m.mu.RUnlock()
namespaces := make([]string, 0, len(m.records))
for name, record := range m.records {
if record.Phase == NamespacePhaseRemoved {
continue
}
namespaces = append(namespaces, name)
}
sort.Strings(namespaces)
return namespaces
}
func (m *inMemoryNamespaceManager) SnapshotRecords() []NamespaceRecord {
m.mu.RLock()
defer m.mu.RUnlock()
records := make([]NamespaceRecord, 0, len(m.records))
for _, record := range m.records {
records = append(records, record.Clone())
}
sort.Slice(records, func(i, j int) bool {
return records[i].Name < records[j].Name
})
return records
}
func (m *inMemoryNamespaceManager) Get(name string) (NamespaceRecord, bool) {
m.mu.RLock()
defer m.mu.RUnlock()
record, ok := m.records[name]
if !ok {
return NamespaceRecord{}, false
}
return record.Clone(), true
}
func (m *inMemoryNamespaceManager) Upsert(event NamespaceEvent) NamespaceRecord {
m.mu.Lock()
defer m.mu.Unlock()
now := event.ObservedAt
if now.IsZero() {
now = time.Now().UTC()
}
record, exists := m.records[event.Name]
if !exists {
record = NamespaceRecord{
Name: event.Name,
RegisteredParts: make(map[string]NamespacePartState),
}
if event.Type == NamespaceEventRemove {
record.Generation = 0
} else {
record.Generation = 1
}
} else if event.Type != NamespaceEventRemove {
record.Generation++
}
record.Name = event.Name
record.Source = event.Source
record.Labels = cloneStringMap(event.Labels)
record.UpdatedAt = now
switch event.Type {
case NamespaceEventAdd, NamespaceEventUpdate, NamespaceEventResync:
record.Phase = NamespacePhaseDiscovered
record.LastError = ""
case NamespaceEventRemove:
record.Phase = NamespacePhaseRemoved
}
if record.RegisteredParts == nil {
record.RegisteredParts = make(map[string]NamespacePartState)
}
m.records[event.Name] = record
return record.Clone()
}
func (m *inMemoryNamespaceManager) MarkPartState(namespace string, part string, state NamespacePartState) (NamespaceRecord, bool) {
m.mu.Lock()
defer m.mu.Unlock()
record, ok := m.records[namespace]
if !ok {
return NamespaceRecord{}, false
}
if record.RegisteredParts == nil {
record.RegisteredParts = make(map[string]NamespacePartState)
}
if state.UpdatedAt.IsZero() {
state.UpdatedAt = time.Now().UTC()
}
record.RegisteredParts[part] = state
record.Phase = deriveNamespacePhase(record)
record.UpdatedAt = state.UpdatedAt
m.records[namespace] = record
return record.Clone(), true
}
func (m *inMemoryNamespaceManager) MarkPartRegistering(namespace string, part string) (NamespaceRecord, bool) {
return m.MarkPartState(namespace, part, NamespacePartState{State: NamespacePartStateRegistering})
}
func (m *inMemoryNamespaceManager) MarkPartActive(namespace string, part string) (NamespaceRecord, bool) {
return m.MarkPartState(namespace, part, NamespacePartState{State: NamespacePartStateActive})
}
func (m *inMemoryNamespaceManager) MarkPartFailed(namespace string, part string, err error) (NamespaceRecord, bool) {
lastError := ""
if err != nil {
lastError = err.Error()
}
return m.MarkPartState(namespace, part, NamespacePartState{State: NamespacePartStateFailed, LastError: lastError})
}
func (m *inMemoryNamespaceManager) Remove(name string) bool {
m.mu.Lock()
defer m.mu.Unlock()
if _, ok := m.records[name]; !ok {
return false
}
delete(m.records, name)
return true
}
func deriveNamespacePhase(record NamespaceRecord) NamespacePhase {
if len(record.RegisteredParts) == 0 {
return record.Phase
}
hasRegistering := false
for _, part := range record.RegisteredParts {
switch part.State {
case NamespacePartStateFailed:
return NamespacePhaseFailed
case NamespacePartStateRegistering:
hasRegistering = true
case NamespacePartStateActive:
continue
default:
hasRegistering = true
}
}
if hasRegistering {
return NamespacePhaseRegistering
}
return NamespacePhaseActive
}