Files
replicator/replicator/map.go
2023-11-05 20:00:22 -07:00

85 lines
1.6 KiB
Go

package replicator
import (
"context"
"sync"
"time"
)
type Map struct {
m map[Key]ValueVersion
lock *sync.RWMutex
}
func NewMap() Map {
return Map{
m: make(map[Key]ValueVersion),
lock: &sync.RWMutex{},
}
}
func (m Map) KeysSince(ctx context.Context, t time.Time) (chan KeyVersion, *error) {
result := make(chan KeyVersion)
var err error
go func() {
defer close(result)
m.lock.RLock()
defer m.lock.RUnlock()
for k := range m.m {
select {
case result <- KeyVersion{Key: k, Version: m.m[k].Version}:
case <-ctx.Done():
err = ctx.Err()
return
}
}
}()
return result, &err
}
func (m Map) Get(_ context.Context, k Key) (ValueVersion, error) {
m.lock.RLock()
defer m.lock.RUnlock()
return m.m[k], nil
}
func (m Map) Set(_ context.Context, key Key, value Value, version Version) error {
m.lock.Lock()
defer m.lock.Unlock()
if version != nil {
if was, ok := m.m[key]; !ok {
} else if wasVersion, err := was.Version.AsTime(); err != nil {
return err
} else if wantVersion, err := version.AsTime(); err != nil {
return err
} else if wantVersion.Before(wasVersion) {
return nil // conflict
}
}
m.m[key] = ValueVersion{Value: value, Version: version}
return nil
}
func (m Map) Del(_ context.Context, k Key, v Version) error {
m.lock.Lock()
defer m.lock.Unlock()
if v != nil {
if was, ok := m.m[k]; !ok {
return nil
} else if wasVersion, err := was.Version.AsTime(); err != nil {
return err
} else if wantVersion, err := v.AsTime(); err != nil {
return err
} else if wantVersion.Before(wasVersion) {
return nil // conflict
}
}
delete(m.m, k)
return nil
}