Skip to content
Merged
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.

31 changes: 28 additions & 3 deletions adapter/admin_grpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (

"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 @@ -137,6 +138,18 @@ type AdminServer struct {
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 @@ -155,6 +168,9 @@ func NewAdminServer(self NodeIdentity, members []NodeIdentity, opts ...AdminOpti
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 {
Expand Down Expand Up @@ -521,7 +537,10 @@ func (s *AdminServer) GetNodeVersion(
context.Context,
*pb.GetNodeVersionRequest,
) (*pb.GetNodeVersionResponse, error) {
return &pb.GetNodeVersionResponse{NodeVersion: s.nodeVersion}, nil
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 {
Expand Down Expand Up @@ -703,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
}
Comment thread
bootjp marked this conversation as resolved.

// 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 @@ -736,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 @@ -750,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
22 changes: 22 additions & 0 deletions adapter/admin_grpc_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1021,6 +1021,28 @@ func TestAdminTokenAuthSkipsOtherServices(t *testing.T) {
}
}

func TestAdminTokenAuthProtectsInternalAdminMethods(t *testing.T) {
t.Parallel()
unary, _ := AdminTokenAuth("s3cret")
handler := func(_ context.Context, _ any) (any, error) { return "ok", nil }
methods := []string{
pb.Internal_ForwardAdminProposal_FullMethodName,
pb.Internal_ForwardLeaseRead_FullMethodName,
}
for _, method := range methods {
t.Run(method, func(t *testing.T) {
t.Parallel()
info := &grpc.UnaryServerInfo{FullMethod: method}
_, err := unary(context.Background(), nil, info, handler)
require.Equal(t, codes.Unauthenticated, status.Code(err))
ctx := metadata.NewIncomingContext(context.Background(), metadata.Pairs("authorization", "Bearer s3cret"))
resp, err := unary(ctx, nil, info, handler)
require.NoError(t, err)
require.Equal(t, "ok", resp)
})
}
}

func TestAdminTokenAuthEmptyTokenDisabled(t *testing.T) {
t.Parallel()
unary, stream := AdminTokenAuth("")
Expand Down
83 changes: 78 additions & 5 deletions adapter/internal.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import (
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
"github.com/cockroachdb/errors"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)

type InternalOption func(*Internal)
Expand All @@ -18,6 +20,18 @@ func WithInternalTimestampAllocator(alloc kv.TimestampAllocator) InternalOption
}
}

func WithInternalAdminProposer(proposer raftengine.Proposer) InternalOption {
return func(i *Internal) {
i.adminProposer = proposer
}
}

func WithInternalLastCommitTimestamp(reader func() uint64) InternalOption {
return func(i *Internal) {
i.lastCommitTimestamp = reader
}
}

func NewInternalWithEngine(txm kv.Transactional, leader raftengine.LeaderView, clock *kv.HLC, relay *RedisPubSubRelay, opts ...InternalOption) *Internal {
i := &Internal{
leader: leader,
Expand All @@ -32,11 +46,13 @@ func NewInternalWithEngine(txm kv.Transactional, leader raftengine.LeaderView, c
}

type Internal struct {
leader raftengine.LeaderView
transactionManager kv.Transactional
clock *kv.HLC
tsAllocator kv.TimestampAllocator
relay *RedisPubSubRelay
leader raftengine.LeaderView
transactionManager kv.Transactional
clock *kv.HLC
tsAllocator kv.TimestampAllocator
adminProposer raftengine.Proposer
lastCommitTimestamp func() uint64
relay *RedisPubSubRelay

pb.UnimplementedInternalServer
}
Expand Down Expand Up @@ -78,6 +94,63 @@ func (i *Internal) Forward(ctx context.Context, req *pb.ForwardRequest) (*pb.For
}, nil
}

func (i *Internal) ForwardAdminProposal(
ctx context.Context,
req *pb.ForwardAdminProposalRequest,
) (*pb.ForwardAdminProposalResponse, error) {
Comment thread
bootjp marked this conversation as resolved.
if i.leader == nil || i.leader.State() != raftengine.StateLeader {
return nil, errors.WithStack(ErrNotLeader)
}
if err := i.leader.VerifyLeader(ctx); err != nil {
return nil, errors.WithStack(ErrNotLeader)
}
if i.adminProposer == nil {
return nil, errors.New("admin proposer is unavailable")
}
result, err := i.adminProposer.ProposeAdmin(ctx, req.GetPayload())
Comment thread
bootjp marked this conversation as resolved.
if err != nil {
return nil, errors.WithStack(err)
}
if err := forwardedAdminProposalResponseError(result); err != nil {
Comment thread
bootjp marked this conversation as resolved.
return nil, err
Comment thread
bootjp marked this conversation as resolved.
}
return &pb.ForwardAdminProposalResponse{CommitIndex: result.CommitIndex}, nil
}

func (i *Internal) ForwardLeaseRead(
ctx context.Context,
_ *pb.ForwardLeaseReadRequest,
) (*pb.ForwardLeaseReadResponse, error) {
if i.leader == nil || i.leader.State() != raftengine.StateLeader {
return nil, errors.WithStack(ErrNotLeader)
}
index, err := i.leader.LinearizableRead(ctx)
if err != nil {
return nil, errors.WithStack(err)
}
lastCommitTS := uint64(0)
if i.lastCommitTimestamp != nil {
lastCommitTS = i.lastCommitTimestamp()
}
return &pb.ForwardLeaseReadResponse{AppliedIndex: index, LastCommitTs: lastCommitTS}, nil
}

func forwardedAdminProposalResponseError(result *raftengine.ProposalResult) error {
if result == nil {
return errors.New("admin proposal returned nil result")
}
if result.Response == nil {
return nil
}
if err, ok := result.Response.(error); ok {
if errors.Is(err, kv.ErrTooManyActiveBackups) {
return status.Errorf(codes.ResourceExhausted, "%s", kv.ErrTooManyActiveBackups)
}
return errors.WithStack(err)
}
return errors.Errorf("unexpected admin proposal response %T", result.Response)
}

func (i *Internal) RelayPublish(_ context.Context, req *pb.RelayPublishRequest) (*pb.RelayPublishResponse, error) {
if req == nil || i.relay == nil {
return &pb.RelayPublishResponse{}, nil
Expand Down
130 changes: 130 additions & 0 deletions adapter/internal_admin_proposal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package adapter

import (
"context"
stderrors "errors"
"testing"

"github.com/bootjp/elastickv/internal/raftengine"
"github.com/bootjp/elastickv/kv"
pb "github.com/bootjp/elastickv/proto"
"github.com/stretchr/testify/require"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)

type internalAdminLeaderView struct {
state raftengine.State
readIndex uint64
}

func (v internalAdminLeaderView) State() raftengine.State { return v.state }
func (internalAdminLeaderView) Leader() raftengine.LeaderInfo {
return raftengine.LeaderInfo{ID: "leader", Address: "leader:50051"}
}
func (internalAdminLeaderView) VerifyLeader(context.Context) error { return nil }
func (v internalAdminLeaderView) LinearizableRead(context.Context) (uint64, error) {
return v.readIndex, nil
}

func TestInternalForwardLeaseReadUsesLeaderBarrier(t *testing.T) {
t.Parallel()
internal := NewInternalWithEngine(
nil,
internalAdminLeaderView{state: raftengine.StateLeader, readIndex: 23},
nil,
nil,
WithInternalLastCommitTimestamp(func() uint64 { return 37 }),
)
resp, err := internal.ForwardLeaseRead(context.Background(), &pb.ForwardLeaseReadRequest{})
require.NoError(t, err)
require.Equal(t, uint64(23), resp.GetAppliedIndex())
require.Equal(t, uint64(37), resp.GetLastCommitTs())

internal.leader = internalAdminLeaderView{state: raftengine.StateFollower}
_, err = internal.ForwardLeaseRead(context.Background(), &pb.ForwardLeaseReadRequest{})
require.ErrorIs(t, err, ErrNotLeader)
}

type internalAdminProposer struct {
payload []byte
response any
}

func (p *internalAdminProposer) Propose(context.Context, []byte) (*raftengine.ProposalResult, error) {
return nil, stderrors.New("unexpected user proposal")
}

func (p *internalAdminProposer) ProposeAdmin(
_ context.Context,
payload []byte,
) (*raftengine.ProposalResult, error) {
p.payload = append([]byte(nil), payload...)
return &raftengine.ProposalResult{CommitIndex: 17, Response: p.response}, nil
}

func TestInternalForwardAdminProposalUsesLeaderProposer(t *testing.T) {
t.Parallel()
proposer := &internalAdminProposer{}
internal := NewInternalWithEngine(
nil,
internalAdminLeaderView{state: raftengine.StateLeader},
nil,
nil,
WithInternalAdminProposer(proposer),
)

resp, err := internal.ForwardAdminProposal(
context.Background(),
&pb.ForwardAdminProposalRequest{Payload: []byte("pin")},
)
require.NoError(t, err)
require.Equal(t, uint64(17), resp.GetCommitIndex())
require.Equal(t, []byte("pin"), proposer.payload)
}

func TestInternalForwardAdminProposalFailsClosed(t *testing.T) {
t.Parallel()
tests := []struct {
name string
state raftengine.State
proposer raftengine.Proposer
errTarget error
errCode codes.Code
errText string
}{
{name: "follower", state: raftengine.StateFollower, errTarget: ErrNotLeader},
{
name: "apply response", state: raftengine.StateLeader,
proposer: &internalAdminProposer{response: stderrors.New("apply failed")},
errText: "apply failed",
},
{
name: "capacity response", state: raftengine.StateLeader,
proposer: &internalAdminProposer{response: kv.ErrTooManyActiveBackups},
errCode: codes.ResourceExhausted,
errText: kv.ErrTooManyActiveBackups.Error(),
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
internal := NewInternalWithEngine(
nil,
internalAdminLeaderView{state: tc.state},
nil,
nil,
WithInternalAdminProposer(tc.proposer),
)
_, err := internal.ForwardAdminProposal(context.Background(), &pb.ForwardAdminProposalRequest{})
if tc.errTarget != nil {
require.ErrorIs(t, err, tc.errTarget)
}
if tc.errCode != codes.OK {
require.Equal(t, tc.errCode, status.Code(err))
}
if tc.errText != "" {
require.ErrorContains(t, err, tc.errText)
}
})
}
}
7 changes: 7 additions & 0 deletions distribution/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,13 @@ func routesFromCatalog(routes []RouteDescriptor) ([]Route, error) {
return out, nil
}

// RoutesFromCatalogSnapshot validates and materializes the immutable route
// view contained in a durable catalog snapshot. Callers that need ownership
// as of an MVCC timestamp must use this view instead of the live Engine.
func RoutesFromCatalogSnapshot(snapshot CatalogSnapshot) ([]Route, error) {
return routesFromCatalog(snapshot.Routes)
}

func validateRouteOrder(routes []Route) error {
if len(routes) < minRouteCountForOrderValidation {
return nil
Expand Down
Loading
Loading