Skip to content
File

Blob: kv/memory/kv.go

go178 lines
1package memory
2 
3import (
4 "context"
5 "strings"
6 "sync/atomic"
7 
8 "go.miragespace.co/specter/spec/chord"
9 "go.miragespace.co/specter/spec/protocol"
10 
11 "github.com/zhangyunhao116/skipmap"
12 "github.com/zhangyunhao116/skipset"
13)
14 
15type MemoryKV struct {
16 s *skipmap.Uint64Map[*skipmap.StringMap[*kvValue]]
17 hashFn chord.HashFn
18}
19 
20var _ chord.KVProvider = (*MemoryKV)(nil)
21 
22type kvValue struct {
23 simple atomic.Pointer[[]byte]
24 lease atomic.Uint64
25 children *skipset.StringSet
26}
27 
28func (v *kvValue) isDeleted() bool {
29 plain := *v.simple.Load()
30 children := v.children.Len()
31 token := v.lease.Load()
32 deleted := (plain == nil) && (children == 0) && (token == 0)
33 return deleted
34}
35 
36func newInnerMapFunc() *skipmap.StringMap[*kvValue] {
37 return skipmap.NewString[*kvValue]()
38}
39 
40func newValueFunc() *kvValue {
41 v := &kvValue{
42 simple: atomic.Pointer[[]byte]{},
43 children: skipset.NewString(),
44 lease: atomic.Uint64{},
45 }
46 var empty []byte
47 v.simple.Store(&empty)
48 v.lease.Store(0)
49 return v
50}
51 
52func WithHashFn(fn chord.HashFn) *MemoryKV {
53 return &MemoryKV{
54 s: skipmap.NewUint64[*skipmap.StringMap[*kvValue]](),
55 hashFn: fn,
56 }
57}
58 
59func (m *MemoryKV) fetchVal(key []byte) (*kvValue, bool) {
60 p := m.hashFn(key)
61 sKey := string(key)
62 
63 kMap, _ := m.s.LoadOrStoreLazy(p, newInnerMapFunc)
64 return kMap.LoadOrStoreLazy(sKey, newValueFunc)
65}
66 
67func (m *MemoryKV) lookupVal(key []byte) (*kvValue, bool) {
68 kMap, ok := m.s.Load(m.hashFn(key))
69 if !ok {
70 return nil, false
71 }
72 return kMap.Load(string(key))
73}
74 
75// delete plain and prefix keyspaces
76func (m *MemoryKV) deleteAll(key []byte) {
77 p := m.hashFn(key)
78 sKey := string(key)
79 
80 if kMap, ok := m.s.Load(p); ok {
81 kMap.Delete(sKey)
82 }
83}
84 
85func (m *MemoryKV) Import(ctx context.Context, keys [][]byte, values []*protocol.KVTransfer) error {
86 for i, key := range keys {
87 v, _ := m.fetchVal(key)
88 bytes := values[i].GetSimpleValue()
89 v.simple.Store(&bytes)
90 v.lease.Store(values[i].GetLeaseToken())
91 for _, child := range values[i].GetPrefixChildren() {
92 v.children.Add(string(child))
93 }
94 }
95 return nil
96}
97 
98func (m *MemoryKV) ListKeys(_ context.Context, prefix []byte) ([]*protocol.KeyComposite, error) {
99 keys := make([]*protocol.KeyComposite, 0)
100 
101 m.s.Range(func(_ uint64, kMap *skipmap.StringMap[*kvValue]) bool {
102 kMap.Range(func(key string, v *kvValue) bool {
103 if !strings.HasPrefix(key, string(prefix)) {
104 return true
105 }
106 if *v.simple.Load() != nil {
107 keys = append(keys, &protocol.KeyComposite{
108 Type: protocol.KeyComposite_SIMPLE,
109 Key: []byte(key),
110 })
111 }
112 if v.children.Len() > 0 {
113 keys = append(keys, &protocol.KeyComposite{
114 Type: protocol.KeyComposite_PREFIX,
115 Key: []byte(key),
116 })
117 }
118 if v.lease.Load() != 0 {
119 keys = append(keys, &protocol.KeyComposite{
120 Type: protocol.KeyComposite_LEASE,
121 Key: []byte(key),
122 })
123 }
124 return true
125 })
126 return true
127 })
128 
129 return keys, nil
130}
131 
132func (m *MemoryKV) Export(_ context.Context, keys [][]byte) ([]*protocol.KVTransfer, error) {
133 vals := make([]*protocol.KVTransfer, len(keys))
134 for i, key := range keys {
135 v, ok := m.lookupVal(key)
136 if !ok {
137 vals[i] = &protocol.KVTransfer{PrefixChildren: make([][]byte, 0)}
138 continue
139 }
140 plain := *v.simple.Load()
141 children, _ := m.PrefixList(context.Background(), key)
142 token := v.lease.Load()
143 vals[i] = &protocol.KVTransfer{
144 SimpleValue: plain,
145 PrefixChildren: children,
146 LeaseToken: token,
147 }
148 }
149 return vals, nil
150}
151 
152func (m *MemoryKV) RangeKeys(ctx context.Context, low, high uint64) ([][]byte, error) {
153 keys := make([][]byte, 0)
154 
155 m.s.Range(func(id uint64, kMap *skipmap.StringMap[*kvValue]) bool {
156 if chord.Between(low, id, high, true) {
157 kMap.Range(func(key string, v *kvValue) bool {
158 if v.isDeleted() {
159 return true
160 }
161 keys = append(keys, []byte(key))
162 return true
163 })
164 }
165 return true
166 })
167 
168 return keys, nil
169}
170 
171func (m *MemoryKV) RemoveKeys(ctx context.Context, keys [][]byte) error {
172 for _, key := range keys {
173 m.deleteAll(key)
174 }
175 
176 return nil
177}