Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,355 changes: 1,355 additions & 0 deletions adapter/admin_backup.go

Large diffs are not rendered by default.

982 changes: 982 additions & 0 deletions adapter/admin_backup_test.go

Large diffs are not rendered by default.

274 changes: 262 additions & 12 deletions adapter/admin_grpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,12 @@ import (
"strconv"
"strings"
"sync"
"sync/atomic"
"time"

"github.com/bootjp/elastickv/internal/raftengine"
"github.com/bootjp/elastickv/keyviz"
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
"github.com/cockroachdb/errors"
"google.golang.org/grpc"
Expand Down Expand Up @@ -41,6 +43,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
Expand All @@ -55,6 +58,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
Expand All @@ -79,6 +131,25 @@ 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

backupMu sync.Mutex
backupStateMu sync.Mutex
backupStore BackupStore
backupReadFence BackupReadFence
backupPeerProbe BackupPeerProbe
backupLimiter BackupPinLimiter
backupTokenKey [32]byte
backupProtocolVersion uint32
backupConfig backupConfig
backupProposers map[uint64]raftengine.Proposer
backupSessions map[kv.BackupPinID]backupSession

pb.UnimplementedAdminServer
}

Expand All @@ -87,15 +158,26 @@ 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,
capabilities: make(map[string]bool),
groups: make(map[uint64]AdminGroup),
now: time.Now,
srv := &AdminServer{
self: self,
members: cloned,
capabilities: make(map[string]bool),
groups: make(map[uint64]AdminGroup),
now: time.Now,
leaderVersionProbeTimeout: defaultAdminLeaderVersionProbeTimeout,
leaderVersionCacheTTL: defaultAdminLeaderVersionCacheTTL,
backupConfig: defaultBackupConfig(),
backupProposers: make(map[uint64]raftengine.Proposer),
backupSessions: make(map[kv.BackupPinID]backupSession),
}
for _, opt := range opts {
if opt != nil {
opt(srv)
}
}
return srv
}

// SetClock overrides the clock used by GetRaftGroups, letting tests inject a
Expand Down Expand Up @@ -410,16 +492,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
Expand All @@ -440,11 +527,168 @@ 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,
BackupProtocolVersion: s.backupProtocolVersion,
}, 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
}

func (s *AdminServer) snapshotLeaders() []*pb.GroupLeader {
s.groupsMu.RLock()
defer s.groupsMu.RUnlock()
Expand Down Expand Up @@ -478,6 +722,12 @@ func sortedGroupIDs(m map[uint64]AdminGroup) []uint64 {
// package-qualify the service name) does not silently bypass the auth gate.
var adminMethodPrefix = "/" + pb.Admin_ServiceDesc.ServiceName + "/"

func adminAuthenticatedMethod(fullMethod string) bool {
return strings.HasPrefix(fullMethod, adminMethodPrefix) ||
fullMethod == pb.Internal_ForwardAdminProposal_FullMethodName ||
fullMethod == pb.Internal_ForwardLeaseRead_FullMethodName
}

// AdminTokenAuth builds a gRPC unary+stream interceptor pair enforcing
// "authorization: Bearer <token>" metadata against the supplied token. An
// empty token disables enforcement; callers should pair that mode with a
Expand Down Expand Up @@ -511,7 +761,7 @@ func AdminTokenAuth(token string) (grpc.UnaryServerInterceptor, grpc.StreamServe
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (any, error) {
if !strings.HasPrefix(info.FullMethod, adminMethodPrefix) {
if !adminAuthenticatedMethod(info.FullMethod) {
return handler(ctx, req)
}
if err := check(ctx); err != nil {
Expand All @@ -525,7 +775,7 @@ func AdminTokenAuth(token string) (grpc.UnaryServerInterceptor, grpc.StreamServe
info *grpc.StreamServerInfo,
handler grpc.StreamHandler,
) error {
if !strings.HasPrefix(info.FullMethod, adminMethodPrefix) {
if !adminAuthenticatedMethod(info.FullMethod) {
return handler(srv, ss)
}
if err := check(ss.Context()); err != nil {
Expand Down
Loading