Skip to content
File

Blob: chord/server_rpc.go

go248 lines
1package chord
2 
3import (
4 "context"
5 
6 "go.miragespace.co/specter/spec/chord"
7 "go.miragespace.co/specter/spec/protocol"
8 "go.miragespace.co/specter/spec/rpc"
9)
10 
11type Server struct {
12 LocalNode chord.VNode
13 Factory RemoteNodeFactory
14}
15 
16type RemoteNodeFactory func(*protocol.Node) (chord.VNode, error)
17 
18var _ protocol.KVService = (*Server)(nil)
19var _ protocol.VNodeService = (*Server)(nil)
20 
21func (r *Server) Identity(_ context.Context, _ *protocol.IdentityRequest) (*protocol.IdentityResponse, error) {
22 return &protocol.IdentityResponse{
23 Identity: r.LocalNode.Identity(),
24 }, nil
25}
26 
27func (r *Server) Ping(_ context.Context, _ *protocol.PingRequest) (*protocol.PingResponse, error) {
28 if err := r.LocalNode.Ping(); err != nil {
29 return nil, rpc.WrapError(err)
30 }
31 return &protocol.PingResponse{}, nil
32}
33 
34func (r *Server) Notify(_ context.Context, req *protocol.NotifyRequest) (*protocol.NotifyResponse, error) {
35 predecessor := req.GetPredecessor()
36 
37 vnode, err := r.Factory(predecessor)
38 if err != nil {
39 return nil, rpc.WrapError(err)
40 }
41 err = r.LocalNode.Notify(vnode)
42 if err != nil {
43 return nil, rpc.WrapError(err)
44 }
45 
46 return &protocol.NotifyResponse{}, nil
47}
48 
49func (r *Server) FindSuccessor(_ context.Context, req *protocol.FindSuccessorRequest) (*protocol.FindSuccessorResponse, error) {
50 key := req.GetKey()
51 vnode, err := r.LocalNode.FindSuccessor(key)
52 if err != nil {
53 return nil, rpc.WrapError(err)
54 }
55 return &protocol.FindSuccessorResponse{
56 Successor: vnode.Identity(),
57 }, nil
58}
59 
60func (r *Server) GetSuccessors(_ context.Context, _ *protocol.GetSuccessorsRequest) (*protocol.GetSuccessorsResponse, error) {
61 vnodes, err := r.LocalNode.GetSuccessors()
62 if err != nil {
63 return nil, rpc.WrapError(err)
64 }
65 identities := make([]*protocol.Node, 0)
66 for _, vnode := range vnodes {
67 if vnode == nil {
68 continue
69 }
70 identities = append(identities, vnode.Identity())
71 }
72 return &protocol.GetSuccessorsResponse{
73 Successors: identities,
74 }, nil
75}
76 
77func (r *Server) GetPredecessor(_ context.Context, _ *protocol.GetPredecessorRequest) (*protocol.GetPredecessorResponse, error) {
78 vnode, err := r.LocalNode.GetPredecessor()
79 if err != nil {
80 return nil, err
81 }
82 var pre *protocol.Node
83 if vnode != nil {
84 pre = vnode.Identity()
85 }
86 return &protocol.GetPredecessorResponse{
87 Predecessor: pre,
88 }, nil
89}
90 
91func (r *Server) RequestToJoin(_ context.Context, req *protocol.RequestToJoinRequest) (*protocol.RequestToJoinResponse, error) {
92 joiner := req.GetJoiner()
93 
94 vnode, err := r.Factory(joiner)
95 if err != nil {
96 return nil, rpc.WrapError(err)
97 }
98 
99 pre, vnodes, err := r.LocalNode.RequestToJoin(vnode)
100 if err != nil {
101 return nil, rpc.WrapError(err)
102 }
103 
104 successors := make([]*protocol.Node, 0)
105 for _, vnode := range vnodes {
106 if vnode == nil {
107 continue
108 }
109 successors = append(successors, vnode.Identity())
110 }
111 
112 return &protocol.RequestToJoinResponse{
113 Predecessor: pre.Identity(),
114 Successors: successors,
115 }, nil
116}
117 
118func (r *Server) FinishJoin(_ context.Context, req *protocol.MembershipConclusionRequest) (*protocol.MembershipConclusionResponse, error) {
119 if err := r.LocalNode.FinishJoin(req.GetStabilize(), req.GetRelease()); err != nil {
120 return nil, rpc.WrapError(err)
121 }
122 return &protocol.MembershipConclusionResponse{}, nil
123}
124 
125func (r *Server) RequestToLeave(_ context.Context, req *protocol.RequestToLeaveRequest) (*protocol.RequestToLeaveResponse, error) {
126 leaver := req.GetLeaver()
127 
128 vnode, err := r.Factory(leaver)
129 if err != nil {
130 return nil, rpc.WrapError(err)
131 }
132 
133 if err := r.LocalNode.RequestToLeave(vnode); err != nil {
134 return nil, rpc.WrapError(err)
135 }
136 return &protocol.RequestToLeaveResponse{}, nil
137}
138 
139func (r *Server) FinishLeave(_ context.Context, req *protocol.MembershipConclusionRequest) (*protocol.MembershipConclusionResponse, error) {
140 if err := r.LocalNode.FinishLeave(req.GetStabilize(), req.GetRelease()); err != nil {
141 return nil, rpc.WrapError(err)
142 }
143 return &protocol.MembershipConclusionResponse{}, nil
144}
145 
146func (r *Server) Put(ctx context.Context, req *protocol.SimpleRequest) (*protocol.SimpleResponse, error) {
147 if err := r.LocalNode.Put(ctx, req.GetKey(), req.GetValue()); err != nil {
148 return nil, rpc.WrapErrorKV(string(req.GetKey()), err)
149 }
150 return &protocol.SimpleResponse{}, nil
151}
152 
153func (r *Server) Get(ctx context.Context, req *protocol.SimpleRequest) (*protocol.SimpleResponse, error) {
154 val, err := r.LocalNode.Get(ctx, req.GetKey())
155 if err != nil {
156 return nil, rpc.WrapErrorKV(string(req.GetKey()), err)
157 }
158 return &protocol.SimpleResponse{
159 Value: val,
160 }, nil
161}
162 
163func (r *Server) Delete(ctx context.Context, req *protocol.SimpleRequest) (*protocol.SimpleResponse, error) {
164 err := r.LocalNode.Delete(ctx, req.GetKey())
165 if err != nil {
166 return nil, rpc.WrapErrorKV(string(req.GetKey()), err)
167 }
168 return &protocol.SimpleResponse{}, nil
169}
170 
171func (r *Server) Append(ctx context.Context, req *protocol.PrefixRequest) (*protocol.PrefixResponse, error) {
172 if err := r.LocalNode.PrefixAppend(ctx, req.GetPrefix(), req.GetChild()); err != nil {
173 return nil, rpc.WrapErrorKV(string(req.GetPrefix()), err)
174 }
175 return &protocol.PrefixResponse{}, nil
176}
177 
178func (r *Server) List(ctx context.Context, req *protocol.PrefixRequest) (*protocol.PrefixResponse, error) {
179 children, err := r.LocalNode.PrefixList(ctx, req.GetPrefix())
180 if err != nil {
181 return nil, rpc.WrapErrorKV(string(req.GetPrefix()), err)
182 }
183 return &protocol.PrefixResponse{
184 Children: children,
185 }, nil
186}
187 
188func (r *Server) Contains(ctx context.Context, req *protocol.PrefixRequest) (*protocol.PrefixResponse, error) {
189 exists, err := r.LocalNode.PrefixContains(ctx, req.GetPrefix(), req.GetChild())
190 if err != nil {
191 return nil, rpc.WrapErrorKV(string(req.GetPrefix()), err)
192 }
193 return &protocol.PrefixResponse{
194 Exists: exists,
195 }, nil
196}
197 
198func (r *Server) Remove(ctx context.Context, req *protocol.PrefixRequest) (*protocol.PrefixResponse, error) {
199 if err := r.LocalNode.PrefixRemove(ctx, req.GetPrefix(), req.GetChild()); err != nil {
200 return nil, rpc.WrapErrorKV(string(req.GetPrefix()), err)
201 }
202 return &protocol.PrefixResponse{}, nil
203}
204 
205func (r *Server) Acquire(ctx context.Context, req *protocol.LeaseRequest) (*protocol.LeaseResponse, error) {
206 token, err := r.LocalNode.Acquire(ctx, req.GetLease(), req.GetTtl().AsDuration())
207 if err != nil {
208 return nil, rpc.WrapErrorKV(string(req.GetLease()), err)
209 }
210 return &protocol.LeaseResponse{
211 Token: token,
212 }, nil
213}
214 
215func (r *Server) Renew(ctx context.Context, req *protocol.LeaseRequest) (*protocol.LeaseResponse, error) {
216 token, err := r.LocalNode.Renew(ctx, req.GetLease(), req.GetTtl().AsDuration(), req.GetPrevToken())
217 if err != nil {
218 return nil, rpc.WrapErrorKV(string(req.GetLease()), err)
219 }
220 return &protocol.LeaseResponse{
221 Token: token,
222 }, nil
223}
224 
225func (r *Server) Release(ctx context.Context, req *protocol.LeaseRequest) (*protocol.LeaseResponse, error) {
226 if err := r.LocalNode.Release(ctx, req.GetLease(), req.GetPrevToken()); err != nil {
227 return nil, rpc.WrapErrorKV(string(req.GetLease()), err)
228 }
229 return &protocol.LeaseResponse{}, nil
230}
231 
232func (r *Server) Import(ctx context.Context, req *protocol.ImportRequest) (*protocol.ImportResponse, error) {
233 if err := r.LocalNode.Import(ctx, req.GetKeys(), req.GetValues()); err != nil {
234 return nil, rpc.WrapError(err)
235 }
236 return &protocol.ImportResponse{}, nil
237}
238 
239func (r *Server) ListKeys(ctx context.Context, req *protocol.ListKeysRequest) (*protocol.ListKeysResponse, error) {
240 keys, err := r.LocalNode.ListKeys(ctx, req.GetPrefix())
241 if err != nil {
242 return nil, rpc.WrapErrorKV(string(req.GetPrefix()), err)
243 }
244 return &protocol.ListKeysResponse{
245 Keys: keys,
246 }, nil
247}