File
Blob: chord/remote.go
| 1 | package chord |
| 2 | |
| 3 | import ( |
| 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 | |
| 17 | const ( |
| 18 | rpcTimeout = timing.ChordRPCTimeout |
| 19 | pingTimeout = timing.ChordPingTimeout |
| 20 | ) |
| 21 | |
| 22 | type RemoteNode struct { |
| 23 | baseContext context.Context |
| 24 | baseLogger *zap.Logger |
| 25 | logger *zap.Logger |
| 26 | identity *protocol.Node |
| 27 | chordClient rpc.ChordClient |
| 28 | } |
| 29 | |
| 30 | var _ chord.VNode = (*RemoteNode)(nil) |
| 31 | |
| 32 | func 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 | |
| 54 | func (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 | |
| 72 | func (n *RemoteNode) ID() uint64 { |
| 73 | return n.identity.GetId() |
| 74 | } |
| 75 | |
| 76 | func (n *RemoteNode) Identity() *protocol.Node { |
| 77 | return n.identity |
| 78 | } |
| 79 | |
| 80 | func (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 | |
| 89 | func (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 | |
| 100 | func (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 | |
| 128 | func (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 | |
| 152 | func (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 | |
| 176 | func (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 | |
| 188 | func (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 | |
| 201 | func (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 | |
| 212 | func (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 | |
| 224 | func (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 | |
| 237 | func (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 | |
| 251 | func (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 | |
| 263 | func (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 | |
| 277 | func (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 | |
| 292 | func (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 | |
| 304 | func (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 | |
| 316 | func (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 | |
| 329 | func (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 | |
| 360 | func (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 | |
| 372 | func (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 | |
| 383 | func (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 | } |