169 lines
3.8 KiB
Go
169 lines
3.8 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)
|
|
Subscribe(subscriber NamespaceSubscriber)
|
|
SnapshotSubscribers() []string
|
|
Upsert(event NamespaceEvent) NamespaceRecord
|
|
MarkPartState(namespace string, part string, state NamespacePartState) (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 (m *inMemoryNamespaceManager) Subscribe(subscriber NamespaceSubscriber) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.subs[subscriber.Name()] = subscriber
|
|
}
|
|
|
|
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.UpdatedAt = state.UpdatedAt
|
|
m.records[namespace] = record
|
|
return record.Clone(), true
|
|
}
|
|
|
|
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
|
|
} |