Skip to content
File

Blob: chord/remote.go

go394 lines
1package chord
2 
3import (
4 "context"
5 "time"
6 
7 "go.miragespace.co/specter/spec/chord"
8 "go.miragespace.co/specter/spec/protocol"
9 "go.miragespace.co/specter/spec/rpc"
10 "go.miragespace.co/specter/timing"
11 
12 "github.com/TheZeroSlave/zapsentry"
13 "go.uber.org/zap"
14 "google.golang.org/protobuf/types/known/durationpb"
15)
16 
17const (
18 rpcTimeout = timing.ChordRPCTimeout
19 pingTimeout = timing.ChordPingTimeout
20)
21 
22type RemoteNode struct {
23 baseContext context.Context
24 baseLogger *zap.Logger
25 logger *zap.Logger
26 identity *protocol.Node
27 chordClient rpc.ChordClient
28}
29 
30var _ chord.VNode = (*RemoteNode)(nil)
31 
32func NewRemoteNode(ctx context.Context, baseLogger *zap.Logger, chordClient rpc.ChordClient, peer *protocol.Node) (*RemoteNode, error) {
33 if peer == nil {
34 return nil, chord.ErrNodeNil
35 }
36 
37 n := &RemoteNode{
38 baseContext: ctx,
39 baseLogger: baseLogger,
40 identity: peer,
41 chordClient: chordClient,
42 }
43 
44 if peer.GetUnknown() {
45 if err := n.getIdentity(n.chordClient, peer); err != nil {
46 return nil, err
47 }
48 }
49 n.logger = baseLogger.With(zapsentry.NewScope()).With(zap.String("component", "remoteNode"), zap.Object("node", n.identity))
50 
51 return n, nil
52}
53 
54func (n *RemoteNode) getIdentity(r rpc.ChordClient, node *protocol.Node) error {
55 if node.GetAddress() == "" {
56 return chord.ErrNodeNil
57 }
58 
59 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, node), rpcTimeout)
60 defer cancel()
61 
62 resp, err := r.Identity(ctx, &protocol.IdentityRequest{})
63 if err != nil {
64 n.baseLogger.Error("error obtaining Identity of RemoteNode", zap.Object("peer", node), zap.Error(err))
65 return err
66 }
67 
68 n.identity = resp.GetIdentity()
69 return nil
70}
71 
72func (n *RemoteNode) ID() uint64 {
73 return n.identity.GetId()
74}
75 
76func (n *RemoteNode) Identity() *protocol.Node {
77 return n.identity
78}
79 
80func (n *RemoteNode) Ping() error {
81 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), pingTimeout)
82 defer cancel()
83 
84 _, err := n.chordClient.Ping(ctx, &protocol.PingRequest{})
85 
86 return chord.ErrorMapper(err)
87}
88 
89func (n *RemoteNode) Notify(predecessor chord.VNode) error {
90 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
91 defer cancel()
92 
93 _, err := n.chordClient.Notify(ctx, &protocol.NotifyRequest{
94 Predecessor: predecessor.Identity(),
95 })
96 
97 return chord.ErrorMapper(err)
98}
99 
100func (n *RemoteNode) FindSuccessor(key uint64) (chord.VNode, error) {
101 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
102 defer cancel()
103 
104 resp, err := n.chordClient.FindSuccessor(ctx, &protocol.FindSuccessorRequest{
105 Key: key,
106 })
107 if err != nil {
108 return nil, chord.ErrorMapper(err)
109 }
110 
111 if resp.GetSuccessor() == nil {
112 return nil, nil
113 }
114 
115 if resp.GetSuccessor().GetId() == n.ID() {
116 return n, nil
117 }
118 
119 succ, err := NewRemoteNode(n.baseContext, n.baseLogger, n.chordClient, resp.GetSuccessor())
120 if err != nil {
121 n.logger.Error("error creating new RemoteNode in FindSuccessor", zap.Object("peer", resp.GetSuccessor()), zap.Error(err))
122 return nil, err
123 }
124 
125 return succ, nil
126}
127 
128func (n *RemoteNode) GetSuccessors() ([]chord.VNode, error) {
129 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
130 defer cancel()
131 
132 resp, err := n.chordClient.GetSuccessors(ctx, &protocol.GetSuccessorsRequest{})
133 if err != nil {
134 return nil, chord.ErrorMapper(err)
135 }
136 
137 succList := resp.GetSuccessors()
138 nodes := make([]chord.VNode, 0, len(succList))
139 
140 for _, succ := range succList {
141 node, err := NewRemoteNode(n.baseContext, n.baseLogger, n.chordClient, succ)
142 if err != nil {
143 n.logger.Error("error creating RemoteNote in GetSuccessors", zap.Object("peer", n.Identity()), zap.Error(err))
144 continue
145 }
146 nodes = append(nodes, node)
147 }
148 
149 return nodes, nil
150}
151 
152func (n *RemoteNode) GetPredecessor() (chord.VNode, error) {
153 ctx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
154 defer cancel()
155 
156 resp, err := n.chordClient.GetPredecessor(ctx, &protocol.GetPredecessorRequest{})
157 if err != nil {
158 return nil, chord.ErrorMapper(err)
159 }
160 
161 if resp.GetPredecessor() == nil {
162 return nil, nil
163 }
164 if resp.GetPredecessor().GetId() == n.ID() {
165 return n, nil
166 }
167 
168 pre, err := NewRemoteNode(n.baseContext, n.baseLogger, n.chordClient, resp.GetPredecessor())
169 if err != nil {
170 return nil, err
171 }
172 
173 return pre, nil
174}
175 
176func (n *RemoteNode) Put(ctx context.Context, key, value []byte) error {
177 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
178 defer cancel()
179 
180 _, err := n.chordClient.Put(reqCtx, &protocol.SimpleRequest{
181 Key: key,
182 Value: value,
183 })
184 
185 return chord.ErrorMapper(err)
186}
187 
188func (n *RemoteNode) Get(ctx context.Context, key []byte) ([]byte, error) {
189 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
190 defer cancel()
191 
192 resp, err := n.chordClient.Get(reqCtx, &protocol.SimpleRequest{
193 Key: key,
194 })
195 if err != nil {
196 return nil, chord.ErrorMapper(err)
197 }
198 return resp.GetValue(), nil
199}
200 
201func (n *RemoteNode) Delete(ctx context.Context, key []byte) error {
202 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
203 defer cancel()
204 
205 _, err := n.chordClient.Delete(reqCtx, &protocol.SimpleRequest{
206 Key: key,
207 })
208 
209 return chord.ErrorMapper(err)
210}
211 
212func (n *RemoteNode) PrefixAppend(ctx context.Context, prefix []byte, child []byte) error {
213 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
214 defer cancel()
215 
216 _, err := n.chordClient.Append(reqCtx, &protocol.PrefixRequest{
217 Prefix: prefix,
218 Child: child,
219 })
220 
221 return chord.ErrorMapper(err)
222}
223 
224func (n *RemoteNode) PrefixList(ctx context.Context, prefix []byte) ([][]byte, error) {
225 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
226 defer cancel()
227 
228 resp, err := n.chordClient.List(reqCtx, &protocol.PrefixRequest{
229 Prefix: prefix,
230 })
231 if err != nil {
232 return nil, chord.ErrorMapper(err)
233 }
234 return resp.GetChildren(), nil
235}
236 
237func (n *RemoteNode) PrefixContains(ctx context.Context, prefix []byte, child []byte) (bool, error) {
238 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
239 defer cancel()
240 
241 resp, err := n.chordClient.Contains(reqCtx, &protocol.PrefixRequest{
242 Prefix: prefix,
243 Child: child,
244 })
245 if err != nil {
246 return false, chord.ErrorMapper(err)
247 }
248 return resp.GetExists(), nil
249}
250 
251func (n *RemoteNode) PrefixRemove(ctx context.Context, prefix []byte, child []byte) error {
252 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
253 defer cancel()
254 
255 _, err := n.chordClient.Remove(reqCtx, &protocol.PrefixRequest{
256 Prefix: prefix,
257 Child: child,
258 })
259 
260 return chord.ErrorMapper(err)
261}
262 
263func (n *RemoteNode) Acquire(ctx context.Context, lease []byte, ttl time.Duration) (uint64, error) {
264 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
265 defer cancel()
266 
267 resp, err := n.chordClient.Acquire(reqCtx, &protocol.LeaseRequest{
268 Lease: lease,
269 Ttl: durationpb.New(ttl),
270 })
271 if err != nil {
272 return 0, chord.ErrorMapper(err)
273 }
274 return resp.GetToken(), nil
275}
276 
277func (n *RemoteNode) Renew(ctx context.Context, lease []byte, ttl time.Duration, prevToken uint64) (newToken uint64, err error) {
278 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
279 defer cancel()
280 
281 resp, err := n.chordClient.Renew(reqCtx, &protocol.LeaseRequest{
282 Lease: lease,
283 Ttl: durationpb.New(ttl),
284 PrevToken: prevToken,
285 })
286 if err != nil {
287 return 0, chord.ErrorMapper(err)
288 }
289 return resp.GetToken(), nil
290}
291 
292func (n *RemoteNode) Release(ctx context.Context, lease []byte, token uint64) error {
293 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
294 defer cancel()
295 
296 _, err := n.chordClient.Release(reqCtx, &protocol.LeaseRequest{
297 Lease: lease,
298 PrevToken: token,
299 })
300 
301 return chord.ErrorMapper(err)
302}
303 
304func (n *RemoteNode) Import(ctx context.Context, keys [][]byte, values []*protocol.KVTransfer) error {
305 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
306 defer cancel()
307 
308 _, err := n.chordClient.Import(reqCtx, &protocol.ImportRequest{
309 Keys: keys,
310 Values: values,
311 })
312 
313 return chord.ErrorMapper(err)
314}
315 
316func (n *RemoteNode) ListKeys(ctx context.Context, prefix []byte) ([]*protocol.KeyComposite, error) {
317 reqCtx, cancel := context.WithTimeout(rpc.WithNode(ctx, n.identity), rpcTimeout)
318 defer cancel()
319 
320 resp, err := n.chordClient.ListKeys(reqCtx, &protocol.ListKeysRequest{
321 Prefix: prefix,
322 })
323 if err != nil {
324 return nil, chord.ErrorMapper(err)
325 }
326 return resp.GetKeys(), nil
327}
328 
329func (n *RemoteNode) RequestToJoin(joiner chord.VNode) (chord.VNode, []chord.VNode, error) {
330 reqCtx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
331 defer cancel()
332 
333 resp, err := n.chordClient.RequestToJoin(reqCtx, &protocol.RequestToJoinRequest{
334 Joiner: joiner.Identity(),
335 })
336 if err != nil {
337 return nil, nil, chord.ErrorMapper(err)
338 }
339 
340 pre := resp.GetPredecessor()
341 succList := resp.GetSuccessors()
342 
343 pp, err := NewRemoteNode(n.baseContext, n.baseLogger, n.chordClient, pre)
344 if err != nil {
345 return nil, nil, err
346 }
347 
348 nodes := make([]chord.VNode, 0, len(succList))
349 for _, succ := range succList {
350 node, err := NewRemoteNode(n.baseContext, n.baseLogger, n.chordClient, succ)
351 if err != nil {
352 continue
353 }
354 nodes = append(nodes, node)
355 }
356 
357 return pp, nodes, nil
358}
359 
360func (n *RemoteNode) FinishJoin(stabilize bool, release bool) error {
361 reqCtx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
362 defer cancel()
363 
364 _, err := n.chordClient.FinishJoin(reqCtx, &protocol.MembershipConclusionRequest{
365 Stabilize: stabilize,
366 Release: release,
367 })
368 
369 return chord.ErrorMapper(err)
370}
371 
372func (n *RemoteNode) RequestToLeave(leaver chord.VNode) error {
373 reqCtx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
374 defer cancel()
375 
376 _, err := n.chordClient.RequestToLeave(reqCtx, &protocol.RequestToLeaveRequest{
377 Leaver: leaver.Identity(),
378 })
379 
380 return chord.ErrorMapper(err)
381}
382 
383func (n *RemoteNode) FinishLeave(stabilize bool, release bool) error {
384 reqCtx, cancel := context.WithTimeout(rpc.WithNode(n.baseContext, n.identity), rpcTimeout)
385 defer cancel()
386 
387 _, err := n.chordClient.FinishLeave(reqCtx, &protocol.MembershipConclusionRequest{
388 Stabilize: stabilize,
389 Release: release,
390 })
391 
392 return chord.ErrorMapper(err)
393}