-
Notifications
You must be signed in to change notification settings - Fork 2
backup: add admin version API scaffolding #1059
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
2c44879
8f78f4b
6ece87b
ab741c2
7aece27
5cb7ae0
259d59b
b9686b1
a0e32f6
fd9ea80
e8c0cb2
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -8,6 +8,7 @@ import ( | |
| "strconv" | ||
| "strings" | ||
| "sync" | ||
| "sync/atomic" | ||
| "time" | ||
|
|
||
| "github.com/bootjp/elastickv/internal/raftengine" | ||
|
|
@@ -41,6 +42,7 @@ type KeyVizSampler interface { | |
| type AdminGroup interface { | ||
| Status() raftengine.Status | ||
| Configuration(ctx context.Context) (raftengine.Configuration, error) | ||
| SnapshotEvery() uint64 | ||
| } | ||
|
|
||
| // NodeIdentity is the value form of the protobuf NodeIdentity message used for | ||
|
|
@@ -55,6 +57,55 @@ func (n NodeIdentity) toProto() *pb.NodeIdentity { | |
| return &pb.NodeIdentity{NodeId: n.NodeID, GrpcAddress: n.GRPCAddress} | ||
| } | ||
|
|
||
| // LeaderVersionProbe fetches the Admin service version for a peer address. | ||
| // Implementations must honor ctx for the 500ms async GetRaftGroups probe | ||
| // budget and any auth metadata copied from the inbound Admin request. | ||
| type LeaderVersionProbe func(ctx context.Context, grpcAddress string) (string, error) | ||
|
|
||
| // AdminOption adjusts optional AdminServer behavior without changing existing | ||
| // test construction call sites. | ||
| type AdminOption func(*AdminServer) | ||
|
|
||
| func WithAdminNodeVersion(version string) AdminOption { | ||
| return func(s *AdminServer) { | ||
| s.nodeVersion = version | ||
| } | ||
| } | ||
|
|
||
| func WithAdminLeaderVersionProbe(probe LeaderVersionProbe) AdminOption { | ||
| return func(s *AdminServer) { | ||
| s.leaderVersionProbe = probe | ||
| } | ||
| } | ||
|
|
||
| func WithAdminLeaderVersionProbeTimeout(timeout time.Duration) AdminOption { | ||
| return func(s *AdminServer) { | ||
| if timeout > 0 { | ||
| s.leaderVersionProbeTimeout = timeout | ||
| } | ||
| } | ||
| } | ||
|
|
||
| func WithAdminLeaderVersionCacheTTL(ttl time.Duration) AdminOption { | ||
| return func(s *AdminServer) { | ||
| if ttl > 0 { | ||
| s.leaderVersionCacheTTL = ttl | ||
| } | ||
| } | ||
| } | ||
|
|
||
| type versionCacheEntry struct { | ||
| version string | ||
| fetchedAt time.Time | ||
| probeID uint64 | ||
| } | ||
|
|
||
| const ( | ||
| defaultAdminLeaderVersionProbeTimeout = 500 * time.Millisecond | ||
| defaultAdminLeaderVersionCacheTTL = 10 * time.Second | ||
| minLeaderVersionProbeAttemptTimeout = 100 * time.Millisecond | ||
| ) | ||
|
|
||
| // AdminServer implements the node-side Admin gRPC service described in | ||
| // docs/admin_ui_key_visualizer_design.md §4 (Layer A). Phase 0 only implements | ||
| // GetClusterOverview and GetRaftGroups; remaining RPCs return Unimplemented so | ||
|
|
@@ -78,6 +129,13 @@ type AdminServer struct { | |
| // pairs atomically with concurrent RPC reads. | ||
| sampler KeyVizSampler | ||
|
|
||
| nodeVersion string | ||
| leaderVersionProbe LeaderVersionProbe | ||
| leaderVersionProbeTimeout time.Duration | ||
| leaderVersionCacheTTL time.Duration | ||
| leaderVersionProbeSeq atomic.Uint64 | ||
| versionCache sync.Map | ||
|
|
||
| pb.UnimplementedAdminServer | ||
| } | ||
|
|
||
|
|
@@ -86,14 +144,22 @@ type AdminServer struct { | |
| // snapshot shipped to the admin binary; callers that already have a membership | ||
| // source may pass nil and let the admin binary's fan-out layer discover peers | ||
| // by other means. | ||
| func NewAdminServer(self NodeIdentity, members []NodeIdentity) *AdminServer { | ||
| func NewAdminServer(self NodeIdentity, members []NodeIdentity, opts ...AdminOption) *AdminServer { | ||
| cloned := append([]NodeIdentity(nil), members...) | ||
| return &AdminServer{ | ||
| self: self, | ||
| members: cloned, | ||
| groups: make(map[uint64]AdminGroup), | ||
| now: time.Now, | ||
| srv := &AdminServer{ | ||
| self: self, | ||
| members: cloned, | ||
| groups: make(map[uint64]AdminGroup), | ||
| now: time.Now, | ||
| leaderVersionProbeTimeout: defaultAdminLeaderVersionProbeTimeout, | ||
| leaderVersionCacheTTL: defaultAdminLeaderVersionCacheTTL, | ||
| } | ||
| for _, opt := range opts { | ||
| if opt != nil { | ||
| opt(srv) | ||
| } | ||
| } | ||
| return srv | ||
| } | ||
|
|
||
| // SetClock overrides the clock used by GetRaftGroups, letting tests inject a | ||
|
|
@@ -380,16 +446,21 @@ func mergeSeedMembers(seeds []NodeIdentity, selfID string, live *liveMembers) { | |
| // GetRaftGroups returns per-group state snapshots. Phase 0 wires commit/applied | ||
| // indices only; per-follower contact and term history land in later phases. | ||
| func (s *AdminServer) GetRaftGroups( | ||
| _ context.Context, | ||
| ctx context.Context, | ||
| _ *pb.GetRaftGroupsRequest, | ||
| ) (*pb.GetRaftGroupsResponse, error) { | ||
| s.groupsMu.RLock() | ||
| defer s.groupsMu.RUnlock() | ||
| ids := sortedGroupIDs(s.groups) | ||
| statuses := make([]raftengine.Status, len(ids)) | ||
| for i, id := range ids { | ||
| statuses[i] = s.groups[id].Status() | ||
| } | ||
| probeAddresses := leaderProbeAddressesByNode(statuses) | ||
| out := make([]*pb.RaftGroupState, 0, len(ids)) | ||
| now := s.now() | ||
| for _, id := range ids { | ||
| st := s.groups[id].Status() | ||
| for i, id := range ids { | ||
| st := statuses[i] | ||
| // Translate LastContact (duration since the last contact with the | ||
| // leader, per raftengine.Status) into an absolute unix-ms so UI | ||
| // clients can diff against their own clock instead of having to | ||
|
|
@@ -410,11 +481,165 @@ func (s *AdminServer) GetRaftGroups( | |
| CommitIndex: st.CommitIndex, | ||
| AppliedIndex: st.AppliedIndex, | ||
| LastContactUnixMs: lastContactUnixMs, | ||
| LeaderNodeVersion: s.leaderNodeVersion(ctx, st.Leader, now, probeAddresses[leaderVersionCacheKey(st.Leader)]), | ||
| }) | ||
| } | ||
| return &pb.GetRaftGroupsResponse{Groups: out}, nil | ||
| } | ||
|
|
||
| func (s *AdminServer) GetNodeVersion( | ||
| context.Context, | ||
| *pb.GetNodeVersionRequest, | ||
| ) (*pb.GetNodeVersionResponse, error) { | ||
| return &pb.GetNodeVersionResponse{NodeVersion: s.nodeVersion}, nil | ||
| } | ||
|
|
||
| func (s *AdminServer) leaderNodeVersion(ctx context.Context, leader raftengine.LeaderInfo, now time.Time, addresses []string) string { | ||
| if version, ok := s.localLeaderVersion(leader); ok { | ||
| return version | ||
| } | ||
| key := leaderVersionCacheKey(leader) | ||
| if version, ok := s.cachedLeaderVersion(key, now); ok { | ||
| return version | ||
| } | ||
| if s.leaderVersionProbe == nil || len(addresses) == 0 { | ||
| return "" | ||
| } | ||
| version, reservation, reserved := s.reserveLeaderVersionProbe(key, now) | ||
| if !reserved { | ||
| return version | ||
| } | ||
| s.probeLeaderVersionAsync(ctx, key, addresses, reservation) | ||
| return "" | ||
| } | ||
|
|
||
| func leaderProbeAddressesByNode(statuses []raftengine.Status) map[string][]string { | ||
| addresses := make(map[string][]string) | ||
| seen := make(map[string]map[string]struct{}) | ||
| for _, st := range statuses { | ||
| key := leaderVersionCacheKey(st.Leader) | ||
| address := strings.TrimSpace(st.Leader.Address) | ||
| if key == "" || address == "" { | ||
| continue | ||
| } | ||
| if seen[key] == nil { | ||
| seen[key] = make(map[string]struct{}) | ||
| } | ||
| if _, ok := seen[key][address]; ok { | ||
| continue | ||
| } | ||
| seen[key][address] = struct{}{} | ||
| addresses[key] = append(addresses[key], address) | ||
| } | ||
| return addresses | ||
| } | ||
|
|
||
| func (s *AdminServer) localLeaderVersion(leader raftengine.LeaderInfo) (string, bool) { | ||
| if leader.ID == "" && leader.Address == "" { | ||
| return "", true | ||
| } | ||
| if leader.ID == s.self.NodeID || (leader.Address != "" && leader.Address == s.self.GRPCAddress) { | ||
| return s.nodeVersion, true | ||
| } | ||
| return "", false | ||
| } | ||
|
|
||
| func leaderVersionCacheKey(leader raftengine.LeaderInfo) string { | ||
| if leader.ID != "" { | ||
| return leader.ID | ||
| } | ||
| return leader.Address | ||
| } | ||
|
|
||
| func (s *AdminServer) cachedLeaderVersion(key string, now time.Time) (string, bool) { | ||
| if key == "" || s.leaderVersionCacheTTL <= 0 { | ||
| return "", false | ||
| } | ||
| actual, ok := s.versionCache.Load(key) | ||
| if !ok { | ||
| return "", false | ||
| } | ||
| entry, ok := actual.(versionCacheEntry) | ||
| if !ok { | ||
| s.versionCache.Delete(key) | ||
| return "", false | ||
| } | ||
| if now.Sub(entry.fetchedAt) > s.leaderVersionCacheTTL { | ||
| s.versionCache.CompareAndDelete(key, entry) | ||
| return "", false | ||
| } | ||
| return entry.version, true | ||
| } | ||
|
|
||
| func (s *AdminServer) reserveLeaderVersionProbe(key string, now time.Time) (string, versionCacheEntry, bool) { | ||
| reservation := versionCacheEntry{ | ||
| fetchedAt: now, | ||
| probeID: s.leaderVersionProbeSeq.Add(1), | ||
| } | ||
| for { | ||
| actual, loaded := s.versionCache.LoadOrStore(key, reservation) | ||
| if !loaded { | ||
| return "", reservation, true | ||
| } | ||
| entry, ok := actual.(versionCacheEntry) | ||
| switch { | ||
| case !ok: | ||
| s.versionCache.Store(key, reservation) | ||
| return "", reservation, true | ||
| case now.Sub(entry.fetchedAt) <= s.leaderVersionCacheTTL: | ||
| return entry.version, versionCacheEntry{}, false | ||
| default: | ||
| if s.versionCache.CompareAndSwap(key, entry, reservation) { | ||
| return "", reservation, true | ||
| } | ||
| } | ||
| } | ||
| } | ||
|
|
||
| func (s *AdminServer) probeLeaderVersionAsync(ctx context.Context, key string, addresses []string, reservation versionCacheEntry) { | ||
| probe := s.leaderVersionProbe | ||
| timeout := s.leaderVersionProbeTimeout | ||
| now := s.now | ||
| candidates := append([]string(nil), addresses...) | ||
| attemptTimeout := leaderVersionProbeAttemptTimeout(timeout, len(candidates)) | ||
| var md metadata.MD | ||
| if incoming, ok := metadata.FromIncomingContext(ctx); ok { | ||
| md = incoming.Copy() | ||
| } | ||
| go func() { | ||
| probeCtx, cancel := context.WithTimeout(context.Background(), timeout) | ||
| defer cancel() | ||
| version := "" | ||
| for _, address := range candidates { | ||
| attemptCtx, attemptCancel := context.WithTimeout(probeCtx, attemptTimeout) | ||
| if md != nil { | ||
| attemptCtx = metadata.NewOutgoingContext(attemptCtx, md) | ||
| } | ||
| probed, err := probe(attemptCtx, address) | ||
| attemptCancel() | ||
| if err == nil { | ||
| version = probed | ||
| break | ||
| } | ||
| if probeCtx.Err() != nil { | ||
| break | ||
| } | ||
| } | ||
| s.versionCache.CompareAndSwap(key, reservation, versionCacheEntry{version: version, fetchedAt: now()}) | ||
| }() | ||
| } | ||
|
|
||
| func leaderVersionProbeAttemptTimeout(total time.Duration, candidates int) time.Duration { | ||
| if total <= 0 || candidates <= 1 { | ||
| return total | ||
| } | ||
| perCandidate := total / time.Duration(candidates) | ||
| if perCandidate < minLeaderVersionProbeAttemptTimeout { | ||
| return min(total, minLeaderVersionProbeAttemptTimeout) | ||
| } | ||
| return perCandidate | ||
|
Comment on lines
+636
to
+640
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
In multi-group clusters where the same leader has many group listener addresses, dividing the 500ms leader probe budget by the number of candidates can shrink each Useful? React with 👍 / 👎.
Owner
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Addressed in signed HEAD e8c0cb2. Per-candidate attempts now have a 100ms minimum unless the total probe budget is smaller; the outer total timeout remains the hard cap. Tests cover large candidate counts, a smaller total budget, and the two-candidate split. |
||
| } | ||
|
|
||
| func (s *AdminServer) snapshotLeaders() []*pb.GroupLeader { | ||
| s.groupsMu.RLock() | ||
| defer s.groupsMu.RUnlock() | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.