/
githubmirror
/
etcd
Обзор
Документация
Войти
/
githubmirror
/
etcd
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
server/etcdserver/v3_server.go
1 290 строк
38 KB
Alessio Attilio
apply: fix data inconsistency in txns by skipping range execution
22 май 2026, 13:14
22 май 2026, 13:14
fbba4f4
Код
Авторство
О чём код?
// Copyright 2015 The etcd Authors // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. package etcdserver import ( "context" "encoding/base64" errorspkg "errors" "math" "strconv" "time" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/trace" "go.uber.org/zap" "golang.org/x/crypto/bcrypt" "google.golang.org/protobuf/proto" pb "go.etcd.io/etcd/api/v3/etcdserverpb" "go.etcd.io/etcd/api/v3/version" "go.etcd.io/etcd/pkg/v3/traceutil" "go.etcd.io/etcd/server/v3/auth" "go.etcd.io/etcd/server/v3/etcdserver/api/membership" apply2 "go.etcd.io/etcd/server/v3/etcdserver/apply" "go.etcd.io/etcd/server/v3/etcdserver/errors" "go.etcd.io/etcd/server/v3/etcdserver/txn" "go.etcd.io/etcd/server/v3/features" "go.etcd.io/etcd/server/v3/lease" "go.etcd.io/etcd/server/v3/lease/leasehttp" "go.etcd.io/etcd/server/v3/storage/mvcc" ) const ( // In the health case, there might be a small gap (10s of entries) between // the applied index and committed index. // However, if the committed entries are very heavy to toApply, the gap might grow. // We should stop accepting new proposals if the gap growing to a certain point. maxGapBetweenApplyAndCommitIndex = 5000 maxNormalGap = maxGapBetweenApplyAndCommitIndex maxPriorityGap = 2 * maxGapBetweenApplyAndCommitIndex traceThreshold = 100 * time.Millisecond // The timeout for the node to catch up its applied index, and is used in // lease related operations, such as LeaseRenew and LeaseTimeToLive. applyTimeout = time.Second ) type RaftKV interface { Range(ctx context.Context, r *pb.RangeRequest) (*pb.RangeResponse, error) RangeStream(r *pb.RangeRequest, rs pb.KV_RangeStreamServer) error Put(ctx context.Context, r *pb.PutRequest) (*pb.PutResponse, error) DeleteRange(ctx context.Context, r *pb.DeleteRangeRequest) (*pb.DeleteRangeResponse, error) Txn(ctx context.Context, r *pb.TxnRequest) (*pb.TxnResponse, error) Compact(ctx context.Context, r *pb.CompactionRequest) (*pb.CompactionResponse, error) } type Lessor interface { // LeaseGrant sends LeaseGrant request to raft and toApply it after committed. LeaseGrant(ctx context.Context, r *pb.LeaseGrantRequest) (*pb.LeaseGrantResponse, error) // LeaseRevoke sends LeaseRevoke request to raft and toApply it after committed. LeaseRevoke(ctx context.Context, r *pb.LeaseRevokeRequest) (*pb.LeaseRevokeResponse, error) // LeaseRenew renews the lease with given ID. The renewed TTL is returned. Or an error // is returned. LeaseRenew(ctx context.Context, id lease.LeaseID) (int64, error) // LeaseTimeToLive retrieves lease information. LeaseTimeToLive(ctx context.Context, r *pb.LeaseTimeToLiveRequest) (*pb.LeaseTimeToLiveResponse, error) // LeaseLeases lists all leases. LeaseLeases(ctx context.Context, r *pb.LeaseLeasesRequest) (*pb.LeaseLeasesResponse, error) } type Authenticator interface { AuthEnable(ctx context.Context, r *pb.AuthEnableRequest) (*pb.AuthEnableResponse, error) AuthDisable(ctx context.Context, r *pb.AuthDisableRequest) (*pb.AuthDisableResponse, error) AuthStatus(ctx context.Context, r *pb.AuthStatusRequest) (*pb.AuthStatusResponse, error) Authenticate(ctx context.Context, r *pb.AuthenticateRequest) (*pb.AuthenticateResponse, error) UserAdd(ctx context.Context, r *pb.AuthUserAddRequest) (*pb.AuthUserAddResponse, error) UserDelete(ctx context.Context, r *pb.AuthUserDeleteRequest) (*pb.AuthUserDeleteResponse, error) UserChangePassword(ctx context.Context, r *pb.AuthUserChangePasswordRequest) (*pb.AuthUserChangePasswordResponse, error) UserGrantRole(ctx context.Context, r *pb.AuthUserGrantRoleRequest) (*pb.AuthUserGrantRoleResponse, error) UserGet(ctx context.Context, r *pb.AuthUserGetRequest) (*pb.AuthUserGetResponse, error) UserRevokeRole(ctx context.Context, r *pb.AuthUserRevokeRoleRequest) (*pb.AuthUserRevokeRoleResponse, error) RoleAdd(ctx context.Context, r *pb.AuthRoleAddRequest) (*pb.AuthRoleAddResponse, error) RoleGrantPermission(ctx context.Context, r *pb.AuthRoleGrantPermissionRequest) (*pb.AuthRoleGrantPermissionResponse, error) RoleGet(ctx context.Context, r *pb.AuthRoleGetRequest) (*pb.AuthRoleGetResponse, error) RoleRevokePermission(ctx context.Context, r *pb.AuthRoleRevokePermissionRequest) (*pb.AuthRoleRevokePermissionResponse, error) RoleDelete(ctx context.Context, r *pb.AuthRoleDeleteRequest) (*pb.AuthRoleDeleteResponse, error) UserList(ctx context.Context, r *pb.AuthUserListRequest) (*pb.AuthUserListResponse, error) RoleList(ctx context.Context, r *pb.AuthRoleListRequest) (*pb.AuthRoleListResponse, error) } func (s *EtcdServer) Range(ctx context.Context, r *pb.RangeRequest) (*pb.RangeResponse, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "range", trace.WithAttributes( attribute.String("range_begin", string(r.GetKey())), attribute.String("range_end", string(r.GetRangeEnd())), attribute.Int64("rev", r.GetRevision()), attribute.Int64("limit", r.GetLimit()), attribute.Bool("count_only", r.GetCountOnly()), attribute.Bool("keys_only", r.GetKeysOnly()), )) defer span.End() ctx, trace := traceutil.EnsureTrace(ctx, s.Logger(), "range", traceutil.Field{Key: "range_begin", Value: string(r.Key)}, traceutil.Field{Key: "range_end", Value: string(r.RangeEnd)}, ) var resp *pb.RangeResponse var err error defer func(start time.Time) { txn.WarnOfExpensiveReadOnlyRangeRequest(s.Logger(), s.Cfg.WarningApplyDuration, start, r, resp, err) if resp != nil { trace.AddField( traceutil.Field{Key: "response_count", Value: len(resp.Kvs)}, traceutil.Field{Key: "response_revision", Value: resp.Header.Revision}, ) } trace.LogIfLong(traceThreshold) success := err == nil requestDurationSec.WithLabelValues("Range", strconv.FormatBool(success)).Observe(time.Since(start).Seconds()) }(time.Now()) if !r.Serializable { err = s.read.LinearizableReadNotify(ctx) trace.Step("agreement among raft nodes before linearized reading") if err != nil { return nil, err } } chk := func(ai *auth.AuthInfo) error { return s.authStore.IsRangePermitted(ai, r.Key, r.RangeEnd) } get := func() { resp, _, err = txn.Range(ctx, s.Logger(), s.KV(), r, true) } if serr := s.doSerialize(ctx, chk, get); serr != nil { err = serr return nil, err } return resp, err } func (s *EtcdServer) RangeStream(r *pb.RangeRequest, rs pb.KV_RangeStreamServer) error { ctx := rs.Context() var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "range_streaming", trace.WithAttributes( attribute.String("range_begin", string(r.GetKey())), attribute.String("range_end", string(r.GetRangeEnd())), attribute.Int64("rev", r.GetRevision()), attribute.Int64("limit", r.GetLimit()), attribute.Bool("count_only", r.GetCountOnly()), )) defer span.End() ctx, trace := traceutil.EnsureTrace(ctx, s.Logger(), "range_streaming", traceutil.Field{Key: "range_begin", Value: string(r.Key)}, traceutil.Field{Key: "range_end", Value: string(r.RangeEnd)}, ) if !r.Serializable { err := s.read.LinearizableReadNotify(ctx) trace.Step("agreement among raft nodes before linearized reading") if err != nil { return err } } chk := func(ai *auth.AuthInfo) error { return s.authStore.IsRangePermitted(ai, r.Key, r.RangeEnd) } var err error get := func() { err = s.rangeStream(ctx, r, rs) } if serr := s.doSerialize(ctx, chk, get); serr != nil { err = serr return err } return err } func (s *EtcdServer) rangeStream(ctx context.Context, r *pb.RangeRequest, rs pb.KV_RangeStreamServer) error { if r.CountOnly { resp, _, err := txn.Range(ctx, s.Logger(), s.KV(), r, false) if err != nil { return err } out := &pb.RangeResponse{ Header: &pb.ResponseHeader{Revision: resp.Header.Revision}, Count: resp.Count, } return rs.Send(&pb.RangeStreamResponse{RangeResponse: out}) } totalLimit := r.Limit if totalLimit == 0 { totalLimit = math.MaxInt64 } r.Limit = initialStreamChunkLimit if r.Limit > totalLimit { r.Limit = totalLimit } count := int64(0) var headerRev int64 for { // gofail: var beforeRangeStreamChunk struct{} resp, _, err := txn.Range(ctx, s.Logger(), s.KV(), r, false) if err != nil { return err } // headerRev should represent the latest store revision at the moment // the server starts handling the client request, and remain stable for // the whole response stream. // // If the client did not explicitly pin a revision (r.Revision == 0), // we pin it here to that same initial latest revision, rather than // advancing it for each subsequent range response in the stream. // As a result, writes committed after the stream begins are not reflected // in later range responses of this stream. if headerRev == 0 { headerRev = resp.Header.Revision if r.Revision == 0 { r.Revision = headerRev } } count += int64(len(resp.Kvs)) var nextKey []byte if resp.More { nextKey = append(resp.Kvs[len(resp.Kvs)-1].Key, '\x00') } out := &pb.RangeResponse{Kvs: resp.Kvs} done := !resp.More || count == totalLimit if done { out.Header = &pb.ResponseHeader{Revision: headerRev} out.More = resp.More out.Count = count if resp.More { remaining, cerr := txn.Count(ctx, s.Logger(), s.KV(), nextKey, r.RangeEnd, r.Revision) if cerr != nil { return cerr } out.Count += remaining } } if err := rs.Send(&pb.RangeStreamResponse{RangeResponse: out}); err != nil { return err } if done { return nil } r.Key = nextKey r.Limit = adjustChunkLimit(r.Limit, proto.Size(resp), int(s.Cfg.MaxRequestBytes)) r.Limit = min(r.Limit, totalLimit-count) } } const initialStreamChunkLimit = 10 // adjustChunkLimit picks the next chunk's Limit so each chunk lands near // the target size: too small wastes round-trips, too large produces // oversized chunks. Doubling/halving only outside [0.5x, 2x] avoids // thrashing when responses sit near the boundary. Always returns >= 1, // since Limit=0 means "unlimited" in txn.Range. func adjustChunkLimit(lastLimit int64, lastSize, targetSize int) int64 { switch { case lastSize < targetSize/2: lastLimit *= 2 case lastSize > targetSize*2: lastLimit /= 2 } if lastLimit == 0 { lastLimit = 1 } return lastLimit } func (s *EtcdServer) Put(ctx context.Context, r *pb.PutRequest) (*pb.PutResponse, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "put", trace.WithAttributes( attribute.String("key", string(r.GetKey())), )) defer span.End() ctx = context.WithValue(ctx, traceutil.StartTimeKey{}, time.Now()) resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{Put: r}) if err != nil { return nil, err } return resp.(*pb.PutResponse), nil } func (s *EtcdServer) DeleteRange(ctx context.Context, r *pb.DeleteRangeRequest) (*pb.DeleteRangeResponse, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "delete_range", trace.WithAttributes( attribute.String("range_begin", string(r.GetKey())), attribute.String("range_end", string(r.GetRangeEnd())), )) defer span.End() resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{DeleteRange: r}) if err != nil { return nil, err } return resp.(*pb.DeleteRangeResponse), nil } func (s *EtcdServer) Txn(ctx context.Context, r *pb.TxnRequest) (*pb.TxnResponse, error) { readOnly := txn.IsTxnReadonly(r) var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "txn", trace.WithAttributes( attribute.String("compare_first_key", firstCompareKey(r.GetCompare())), attribute.String("success_first_key", firstOpKey(r.GetSuccess())), attribute.String("success_first_type", firstOpType(r.GetSuccess())), attribute.Int64("success_first_lease", firstOpLease(r.GetSuccess())), attribute.Int("compare_len", len(r.GetCompare())), attribute.Int("success_len", len(r.GetSuccess())), attribute.Int("failure_len", len(r.GetFailure())), attribute.Bool("read_only", readOnly), )) defer span.End() ctx, trace := traceutil.EnsureTrace(ctx, s.Logger(), "transaction", traceutil.Field{Key: "read_only", Value: readOnly}, ) if readOnly { if !txn.IsTxnSerializable(r) { err := s.read.LinearizableReadNotify(ctx) trace.Step("agreement among raft nodes before linearized reading") if err != nil { return nil, err } } var resp *pb.TxnResponse var err error chk := func(ai *auth.AuthInfo) error { return apply2.CheckTxnAuth(s.authStore, ai, s.lessor, r) } defer func(start time.Time) { txn.WarnOfExpensiveReadOnlyTxnRequest(s.Logger(), s.Cfg.WarningApplyDuration, start, r, resp, err) trace.LogIfLong(traceThreshold) success := err == nil requestDurationSec.WithLabelValues("ReadonlyTxn", strconv.FormatBool(success)).Observe(time.Since(start).Seconds()) }(time.Now()) get := func() { resp, _, err = txn.Txn(ctx, s.Logger(), r, s.Cfg.ServerFeatureGate.Enabled(features.TxnModeWriteWithSharedBuffer), s.KV(), s.lessor, false) } if serr := s.doSerialize(ctx, chk, get); serr != nil { return nil, serr } return resp, err } ctx = context.WithValue(ctx, traceutil.StartTimeKey{}, time.Now()) resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{Txn: r}) if err != nil { return nil, err } return resp.(*pb.TxnResponse), nil } func (s *EtcdServer) Compact(ctx context.Context, r *pb.CompactionRequest) (*pb.CompactionResponse, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "compact", trace.WithAttributes( attribute.Bool("is_physical", r.GetPhysical()), attribute.Int64("rev", r.GetRevision()), )) defer span.End() startTime := time.Now() ctx, trace := traceutil.EnsureTrace(ctx, s.Logger(), "compact") result, err := s.processInternalRaftRequestOnce(ctx, &pb.InternalRaftRequest{Compaction: r}) if result != nil && result.Trace != nil { trace = result.Trace defer func() { trace.LogIfLong(traceThreshold) }() applyStart := result.Trace.GetStartTime() result.Trace.SetStartTime(startTime) trace.InsertStep(0, applyStart, "process raft request") } if r.Physical && result != nil && result.Physc != nil { <-result.Physc // The compaction is done deleting keys; the hash is now settled // but the data is not necessarily committed. If there's a crash, // the hash may revert to a hash prior to compaction completing // if the compaction resumes. Force the finished compaction to // commit so it won't resume following a crash. // // `applySnapshot` sets a new backend instance, so we need to acquire the bemu lock. s.bemu.RLock() s.be.ForceCommit() s.bemu.RUnlock() trace.Step("physically toApply compaction") } if err != nil { return nil, err } if result.Err != nil { return nil, result.Err } resp := result.Resp.(*pb.CompactionResponse) if resp == nil { resp = &pb.CompactionResponse{} } if resp.Header == nil { resp.Header = &pb.ResponseHeader{} } resp.Header.Revision = s.kv.Rev() trace.AddField(traceutil.Field{Key: "response_revision", Value: resp.Header.Revision}) return resp, nil } func (s *EtcdServer) LeaseGrant(ctx context.Context, r *pb.LeaseGrantRequest) (*pb.LeaseGrantResponse, error) { // no id given? choose one for r.ID == int64(lease.NoLease) { // only use positive int64 id's r.ID = int64(s.reqIDGen.Next() & ((1 << 63) - 1)) } var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "lease_grant", trace.WithAttributes( attribute.Int64("id", r.ID), attribute.Int64("ttl", r.GetTTL()), )) defer span.End() if err := s.requireAuthInfo(ctx); err != nil { return nil, err } resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{LeaseGrant: r}) if err != nil { return nil, err } return resp.(*pb.LeaseGrantResponse), nil } func (s *EtcdServer) waitAppliedIndex() error { select { case <-s.ApplyWaitCommit(): case <-s.stopping: return errors.ErrStopped case <-time.After(applyTimeout): return errors.ErrTimeoutWaitAppliedIndex } return nil } func (s *EtcdServer) LeaseRevoke(ctx context.Context, r *pb.LeaseRevokeRequest) (*pb.LeaseRevokeResponse, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "lease_revoke", trace.WithAttributes( attribute.Int64("id", r.GetID()), )) defer span.End() if err := s.requireAuthInfo(ctx); err != nil { return nil, err } resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{LeaseRevoke: r}) if err != nil { return nil, err } return resp.(*pb.LeaseRevokeResponse), nil } func (s *EtcdServer) LeaseRenew(ctx context.Context, id lease.LeaseID) (int64, error) { var span trace.Span ctx, span = traceutil.Tracer.Start(ctx, "lease_renew", trace.WithAttributes( attribute.Int64("id", int64(id)), )) defer span.End() if s.isLeader() { // If s.isLeader() returns true, but we fail to ensure the current // member's leadership, there are a couple of possibilities: // 1. current member gets stuck on writing WAL entries; // 2. current member is in network isolation status; // 3. current member isn't a leader anymore (possibly due to #1 above). // In such case, we just return error to client, so that the client can // switch to another member to continue the lease keep-alive operation. if !s.ensureLeadership() { return -1, lease.ErrNotPrimary } // This change aims to make lease renewal faster under high server load // while preserving correctness. If a lease is not found, it might still be in // the process of being created. We must wait for the applied index to advance // to verify whether the lease truly does not exist. if s.FeatureEnabled(features.FastLeaseKeepAlive) { le := s.lessor.Lookup(id) if le == nil { if err := s.waitAppliedIndex(); err != nil { return 0, err } } } else { if err := s.waitAppliedIndex(); err != nil { return 0, err } } if err := s.checkLeaseRenew(ctx, id); err != nil { return 0, err } ttl, err := s.lessor.Renew(id) if err == nil { // already requested to primary lessor(leader) return ttl, nil } if !errorspkg.Is(err, lease.ErrNotPrimary) { return -1, err } } cctx, cancel := context.WithTimeout(ctx, s.Cfg.ReqTimeout()) defer cancel() // renewals don't go through raft; forward to leader manually for cctx.Err() == nil { leader, lerr := s.waitLeader(cctx) if lerr != nil { return -1, lerr } if err := s.checkLeaseRenew(ctx, id); err != nil { return 0, err } for _, url := range leader.PeerURLs { lurl := url + leasehttp.LeasePrefix ttl, err := leasehttp.RenewHTTP(cctx, id, lurl, s.peerRt) if err == nil || errorspkg.Is(err, lease.ErrLeaseNotFound) { return ttl, err } } // Throttle in case of e.g. connection problems. time.Sleep(50 * time.Millisecond) } err := cctx.Err() switch { case errorspkg.Is(err, context.DeadlineExceeded): return -1, errors.ErrTimeout case errorspkg.Is(err, context.Canceled): return -1, errors.ErrCanceled default: s.Logger().Warn("Unexpected lease renew context error", zap.Error(err)) return -1, errors.ErrCanceled } } func (s *EtcdServer) checkLeaseRenew(ctx context.Context, leaseID lease.LeaseID) error { rev := s.AuthStore().Revision() if !s.AuthStore().IsAuthEnabled() { return nil } authInfo, err := s.AuthInfoFromCtx(ctx) if err != nil { return err } if authInfo == nil { return auth.ErrUserEmpty } if s.AuthStore().IsAdminPermitted(authInfo) == nil { return nil } l := s.lessor.Lookup(leaseID) if l != nil { for _, key := range l.Keys() { if err := s.AuthStore().IsPutPermitted(authInfo, []byte(key)); err != nil { return err } } } if rev != s.AuthStore().Revision() { return auth.ErrAuthOldRevision } return nil } func (s *EtcdServer) checkLeaseTimeToLive(ctx context.Context, leaseID lease.LeaseID) (uint64, error) { rev := s.AuthStore().Revision() if !s.AuthStore().IsAuthEnabled() { return rev, nil } authInfo, err := s.AuthInfoFromCtx(ctx) if err != nil { return rev, err } if authInfo == nil { return rev, auth.ErrUserEmpty } if s.AuthStore().IsAdminPermitted(authInfo) == nil { return rev, nil } l := s.lessor.Lookup(leaseID) if l != nil { for _, key := range l.Keys() { if err := s.AuthStore().IsRangePermitted(authInfo, []byte(key), []byte{}); err != nil { return 0, err } } } return rev, nil } func (s *EtcdServer) leaseTimeToLive(ctx context.Context, r *pb.LeaseTimeToLiveRequest) (*pb.LeaseTimeToLiveResponse, error) { if s.isLeader() { if err := s.waitAppliedIndex(); err != nil { return nil, err } // gofail: var beforeLookupWhenLeaseTimeToLive struct{} // primary; timetolive directly from leader le := s.lessor.Lookup(lease.LeaseID(r.ID)) if le == nil { return nil, lease.ErrLeaseNotFound } // TODO: fill out ResponseHeader resp := &pb.LeaseTimeToLiveResponse{Header: &pb.ResponseHeader{}, ID: r.ID, TTL: int64(le.Remaining().Seconds()), GrantedTTL: le.TTL()} if r.Keys { ks := le.Keys() kbs := make([][]byte, len(ks)) for i := range ks { kbs[i] = []byte(ks[i]) } resp.Keys = kbs } // The leasor could be demoted if leader changed during lookup. // We should return error to force retry instead of returning // incorrect remaining TTL. if le.Demoted() { // NOTE: lease.ErrNotPrimary is not retryable error for // client. Instead, uses ErrLeaderChanged. return nil, errors.ErrLeaderChanged } return resp, nil } cctx, cancel := context.WithTimeout(ctx, s.Cfg.ReqTimeout()) defer cancel() // forward to leader for cctx.Err() == nil { leader, err := s.waitLeader(cctx) if err != nil { return nil, err } for _, url := range leader.PeerURLs { lurl := url + leasehttp.LeaseInternalPrefix resp, err := leasehttp.TimeToLiveHTTP(cctx, lease.LeaseID(r.ID), r.Keys, lurl, s.peerRt) if err == nil { return resp.LeaseTimeToLiveResponse, nil } if errorspkg.Is(err, lease.ErrLeaseNotFound) { return nil, err } } } if errorspkg.Is(cctx.Err(), context.DeadlineExceeded) { return nil, errors.ErrTimeout } return nil, errors.ErrCanceled } func (s *EtcdServer) LeaseTimeToLive(ctx context.Context, r *pb.LeaseTimeToLiveRequest) (*pb.LeaseTimeToLiveResponse, error) { if err := s.requireAuthInfo(ctx); err != nil { return nil, err } var rev uint64 var err error if r.Keys { // check RBAC permission only if Keys is true rev, err = s.checkLeaseTimeToLive(ctx, lease.LeaseID(r.ID)) if err != nil { return nil, err } } resp, err := s.leaseTimeToLive(ctx, r) if err != nil { return nil, err } if r.Keys { if s.AuthStore().IsAuthEnabled() && rev != s.AuthStore().Revision() { return nil, auth.ErrAuthOldRevision } } return resp, nil } func (s *EtcdServer) newHeader() *pb.ResponseHeader { return &pb.ResponseHeader{ ClusterId: uint64(s.cluster.ID()), MemberId: uint64(s.MemberID()), Revision: s.KV().Rev(), RaftTerm: s.Term(), } } // LeaseLeases is really ListLeases !??? func (s *EtcdServer) LeaseLeases(ctx context.Context, _ *pb.LeaseLeasesRequest) (*pb.LeaseLeasesResponse, error) { ls := s.lessor.Leases() if err := s.checkLeaseLeases(ctx, ls); err != nil { return nil, err } lss := make([]*pb.LeaseStatus, len(ls)) for i := range ls { lss[i] = &pb.LeaseStatus{ID: int64(ls[i].ID)} } return &pb.LeaseLeasesResponse{Header: s.newHeader(), Leases: lss}, nil } func (s *EtcdServer) checkLeaseLeases(ctx context.Context, leases []*lease.Lease) error { rev := s.AuthStore().Revision() if !s.AuthStore().IsAuthEnabled() { return nil } authInfo, err := s.AuthInfoFromCtx(ctx) if err != nil { return err } if authInfo == nil { return auth.ErrUserEmpty } if err := s.AuthStore().IsAdminPermitted(authInfo); err == nil { return nil } for _, l := range leases { for _, key := range l.Keys() { if err := s.AuthStore().IsRangePermitted(authInfo, []byte(key), []byte{}); err != nil { return err } } } if rev != s.AuthStore().Revision() { return auth.ErrAuthOldRevision } return nil } func (s *EtcdServer) waitLeader(ctx context.Context) (*membership.Member, error) { leader := s.cluster.Member(s.Leader()) for leader == nil { // wait an election dur := time.Duration(s.Cfg.ElectionTicks) * time.Duration(s.Cfg.TickMs) * time.Millisecond select { case <-time.After(dur): leader = s.cluster.Member(s.Leader()) case <-s.stopping: return nil, errors.ErrStopped case <-ctx.Done(): return nil, errors.ErrNoLeader } } if len(leader.PeerURLs) == 0 { return nil, errors.ErrNoLeader } return leader, nil } func (s *EtcdServer) Alarm(ctx context.Context, r *pb.AlarmRequest) (*pb.AlarmResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{Alarm: r}) if err != nil { return nil, err } return resp.(*pb.AlarmResponse), nil } func (s *EtcdServer) AuthEnable(ctx context.Context, r *pb.AuthEnableRequest) (*pb.AuthEnableResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthEnable: r}) if err != nil { return nil, err } return resp.(*pb.AuthEnableResponse), nil } func (s *EtcdServer) AuthDisable(ctx context.Context, r *pb.AuthDisableRequest) (*pb.AuthDisableResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthDisable: r}) if err != nil { return nil, err } return resp.(*pb.AuthDisableResponse), nil } func (s *EtcdServer) AuthStatus(ctx context.Context, r *pb.AuthStatusRequest) (*pb.AuthStatusResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthStatus: r}) if err != nil { return nil, err } return resp.(*pb.AuthStatusResponse), nil } func (s *EtcdServer) Authenticate(ctx context.Context, r *pb.AuthenticateRequest) (*pb.AuthenticateResponse, error) { if err := s.read.LinearizableReadNotify(ctx); err != nil { return nil, err } lg := s.Logger() // fix https://nvd.nist.gov/vuln/detail/CVE-2021-28235 defer func() { if r != nil { r.Password = "" } }() var resp proto.Message for { checkedRevision, err := s.AuthStore().CheckPassword(r.Name, r.Password) if err != nil { if !errorspkg.Is(err, auth.ErrAuthNotEnabled) { lg.Warn( "invalid authentication was requested", zap.String("user", r.Name), zap.Error(err), ) } return nil, err } st, err := s.AuthStore().GenTokenPrefix() if err != nil { return nil, err } // internalReq doesn't need to have Password because the above s.AuthStore().CheckPassword() already did it. // In addition, it will let a WAL entry not record password as a plain text. internalReq := &pb.InternalAuthenticateRequest{ Name: r.Name, SimpleToken: st, } resp, err = s.raftRequest(ctx, &pb.InternalRaftRequest{Authenticate: internalReq}) if err != nil { return nil, err } if checkedRevision == s.AuthStore().Revision() { break } lg.Info("revision when password checked became stale; retrying") } return resp.(*pb.AuthenticateResponse), nil } func (s *EtcdServer) UserAdd(ctx context.Context, r *pb.AuthUserAddRequest) (*pb.AuthUserAddResponse, error) { if r.Options == nil || !r.Options.NoPassword { hashedPassword, err := bcrypt.GenerateFromPassword([]byte(r.Password), s.authStore.BcryptCost()) if err != nil { return nil, err } r.HashedPassword = base64.StdEncoding.EncodeToString(hashedPassword) r.Password = "" } resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserAdd: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserAddResponse), nil } func (s *EtcdServer) UserDelete(ctx context.Context, r *pb.AuthUserDeleteRequest) (*pb.AuthUserDeleteResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserDelete: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserDeleteResponse), nil } func (s *EtcdServer) UserChangePassword(ctx context.Context, r *pb.AuthUserChangePasswordRequest) (*pb.AuthUserChangePasswordResponse, error) { if r.Password != "" { hashedPassword, err := bcrypt.GenerateFromPassword([]byte(r.Password), s.authStore.BcryptCost()) if err != nil { return nil, err } r.HashedPassword = base64.StdEncoding.EncodeToString(hashedPassword) r.Password = "" } resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserChangePassword: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserChangePasswordResponse), nil } func (s *EtcdServer) UserGrantRole(ctx context.Context, r *pb.AuthUserGrantRoleRequest) (*pb.AuthUserGrantRoleResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserGrantRole: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserGrantRoleResponse), nil } func (s *EtcdServer) UserGet(ctx context.Context, r *pb.AuthUserGetRequest) (*pb.AuthUserGetResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserGet: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserGetResponse), nil } func (s *EtcdServer) UserList(ctx context.Context, r *pb.AuthUserListRequest) (*pb.AuthUserListResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserList: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserListResponse), nil } func (s *EtcdServer) UserRevokeRole(ctx context.Context, r *pb.AuthUserRevokeRoleRequest) (*pb.AuthUserRevokeRoleResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthUserRevokeRole: r}) if err != nil { return nil, err } return resp.(*pb.AuthUserRevokeRoleResponse), nil } func (s *EtcdServer) RoleAdd(ctx context.Context, r *pb.AuthRoleAddRequest) (*pb.AuthRoleAddResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleAdd: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleAddResponse), nil } func (s *EtcdServer) RoleGrantPermission(ctx context.Context, r *pb.AuthRoleGrantPermissionRequest) (*pb.AuthRoleGrantPermissionResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleGrantPermission: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleGrantPermissionResponse), nil } func (s *EtcdServer) RoleGet(ctx context.Context, r *pb.AuthRoleGetRequest) (*pb.AuthRoleGetResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleGet: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleGetResponse), nil } func (s *EtcdServer) RoleList(ctx context.Context, r *pb.AuthRoleListRequest) (*pb.AuthRoleListResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleList: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleListResponse), nil } func (s *EtcdServer) RoleRevokePermission(ctx context.Context, r *pb.AuthRoleRevokePermissionRequest) (*pb.AuthRoleRevokePermissionResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleRevokePermission: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleRevokePermissionResponse), nil } func (s *EtcdServer) RoleDelete(ctx context.Context, r *pb.AuthRoleDeleteRequest) (*pb.AuthRoleDeleteResponse, error) { resp, err := s.raftRequest(ctx, &pb.InternalRaftRequest{AuthRoleDelete: r}) if err != nil { return nil, err } return resp.(*pb.AuthRoleDeleteResponse), nil } func (s *EtcdServer) raftRequest(ctx context.Context, r *pb.InternalRaftRequest) (proto.Message, error) { result, err := s.processInternalRaftRequestOnce(ctx, r) if err != nil { trace.SpanFromContext(ctx).RecordError(err) return nil, err } if result.Err != nil { return nil, result.Err } if startTime, ok := ctx.Value(traceutil.StartTimeKey{}).(time.Time); ok && result.Trace != nil { applyStart := result.Trace.GetStartTime() // The trace object is created in toApply. Here reset the start time to trace // the raft request time by the difference between the request start time // and toApply start time result.Trace.SetStartTime(startTime) result.Trace.InsertStep(0, applyStart, "process raft request") result.Trace.LogIfLong(traceThreshold) } return result.Resp, nil } // doSerialize handles the auth logic, with permissions checked by "chk", for a serialized request "get". Returns a non-nil error on authentication failure. func (s *EtcdServer) doSerialize(ctx context.Context, chk func(*auth.AuthInfo) error, get func()) error { trace := traceutil.Get(ctx) ai, err := s.AuthInfoFromCtx(ctx) if err != nil { return err } if ai == nil { // chk expects non-nil AuthInfo; use empty credentials ai = &auth.AuthInfo{} } if err = chk(ai); err != nil { return err } trace.Step("get authentication metadata") // fetch response for serialized request get() // check for stale token revision in case the auth store was updated while // the request has been handled. if ai.Revision != 0 && ai.Revision != s.authStore.Revision() { return auth.ErrAuthOldRevision } return nil } func (s *EtcdServer) processInternalRaftRequestOnce(ctx context.Context, r *pb.InternalRaftRequest) (*apply2.Result, error) { ai := s.getAppliedIndex() ci := s.getCommittedIndex() if exceedsRequestLimit(ai, ci, r, s.FeatureEnabled(features.PriorityRequest)) { return nil, errors.ErrTooManyRequests } r.Header = &pb.RequestHeader{ ID: s.reqIDGen.Next(), } // check authinfo if it is not InternalAuthenticateRequest if r.Authenticate == nil { authInfo, err := s.AuthInfoFromCtx(ctx) if err != nil { return nil, err } if authInfo != nil { r.Header.Username = authInfo.Username r.Header.AuthRevision = authInfo.Revision } } var ( data []byte err error start = time.Now() reqType = getRequestType(r) ) defer func() { success := err == nil requestDurationSec.WithLabelValues(reqType, strconv.FormatBool(success)).Observe(time.Since(start).Seconds()) }() data, err = proto.Marshal(r) if err != nil { return nil, err } if len(data) > int(s.Cfg.MaxRequestBytes) { return nil, errors.ErrRequestTooLarge } id := r.ID if id == 0 { id = r.Header.ID } ch := s.w.Register(id) cctx, cancel := context.WithTimeout(ctx, s.Cfg.ReqTimeout()) defer cancel() span := trace.SpanFromContext(ctx) span.AddEvent("Send raft proposal") err = s.r.Propose(cctx, data) if err != nil { proposalsFailed.Inc() s.w.Trigger(id, nil) // GC wait return nil, err } proposalsPending.Inc() defer proposalsPending.Dec() select { case x := <-ch: span.AddEvent("Receive raft result") return x.(*apply2.Result), nil case <-cctx.Done(): proposalsFailed.Inc() s.w.Trigger(id, nil) // GC wait return nil, s.parseProposeCtxErr(cctx.Err(), start) case <-s.done: return nil, errors.ErrStopped } } func getRequestType(r *pb.InternalRaftRequest) string { switch { case r.Range != nil: return "Range" case r.Put != nil: return "Put" case r.DeleteRange != nil: return "DeleteRange" case r.Txn != nil: return "Txn" case r.Compaction != nil: return "Compaction" case r.LeaseGrant != nil: return "LeaseGrant" case r.LeaseRevoke != nil: return "LeaseRevoke" case r.LeaseCheckpoint != nil: return "LeaseCheckpoint" case r.Alarm != nil: return "Alarm" case r.Authenticate != nil: return "Authenticate" case r.AuthEnable != nil: return "AuthEnable" case r.AuthDisable != nil: return "AuthDisable" case r.AuthStatus != nil: return "AuthStatus" case r.AuthUserAdd != nil: return "AuthUserAdd" case r.AuthUserDelete != nil: return "AuthUserDelete" case r.AuthUserChangePassword != nil: return "AuthUserChangePassword" case r.AuthUserGrantRole != nil: return "AuthUserGrantRole" case r.AuthUserGet != nil: return "AuthUserGet" case r.AuthUserRevokeRole != nil: return "AuthUserRevokeRole" case r.AuthRoleAdd != nil: return "AuthRoleAdd" case r.AuthRoleGrantPermission != nil: return "AuthRoleGrantPermission" case r.AuthRoleGet != nil: return "AuthRoleGet" case r.AuthRoleRevokePermission != nil: return "AuthRoleRevokePermission" case r.AuthRoleDelete != nil: return "AuthRoleDelete" case r.AuthUserList != nil: return "AuthUserList" case r.AuthRoleList != nil: return "AuthRoleList" case r.ClusterVersionSet != nil: return "ClusterVersionSet" case r.ClusterMemberAttrSet != nil: return "ClusterMemberAttrSet" case r.DowngradeInfoSet != nil: return "DowngradeInfoSet" case r.DowngradeVersionTest != nil: return "DowngradeVersionTest" default: return "Unknown" } } // Watchable returns a watchable interface attached to the etcdserver. func (s *EtcdServer) Watchable() mvcc.WatchableKV { return s.KV() } func (s *EtcdServer) AuthInfoFromCtx(ctx context.Context) (*auth.AuthInfo, error) { authInfo, err := s.AuthStore().AuthInfoFromCtx(ctx) if authInfo != nil || err != nil { return authInfo, err } if !s.Cfg.ClientCertAuthEnabled { return nil, nil } authInfo = s.AuthStore().AuthInfoFromTLS(ctx) return authInfo, nil } func (s *EtcdServer) Downgrade(ctx context.Context, r *pb.DowngradeRequest) (*pb.DowngradeResponse, error) { switch r.Action { case pb.DowngradeRequest_VALIDATE: return s.downgradeValidate(ctx, r.Version) case pb.DowngradeRequest_ENABLE: return s.downgradeEnable(ctx, r) case pb.DowngradeRequest_CANCEL: return s.downgradeCancel(ctx) default: return nil, errors.ErrUnknownMethod } } func (s *EtcdServer) downgradeValidate(ctx context.Context, v string) (*pb.DowngradeResponse, error) { resp := &pb.DowngradeResponse{} targetVersion, err := convertToClusterVersion(v) if err != nil { return nil, err } cv := s.ClusterVersion() if cv == nil { return nil, errors.ErrClusterVersionUnavailable } resp.Version = version.Cluster(cv.String()) err = s.Version().DowngradeValidate(ctx, targetVersion) if err != nil { return nil, err } return resp, nil } func (s *EtcdServer) downgradeEnable(ctx context.Context, r *pb.DowngradeRequest) (*pb.DowngradeResponse, error) { lg := s.Logger() targetVersion, err := convertToClusterVersion(r.Version) if err != nil { lg.Warn("reject downgrade request", zap.Error(err)) return nil, err } err = s.Version().DowngradeEnable(ctx, targetVersion) if err != nil { lg.Warn("reject downgrade request", zap.Error(err)) return nil, err } resp := pb.DowngradeResponse{Version: version.Cluster(s.ClusterVersion().String())} return &resp, nil } func (s *EtcdServer) downgradeCancel(ctx context.Context) (*pb.DowngradeResponse, error) { err := s.Version().DowngradeCancel(ctx) if err != nil { s.lg.Warn("failed to cancel downgrade", zap.Error(err)) } resp := pb.DowngradeResponse{Version: version.Cluster(s.ClusterVersion().String())} return &resp, nil } func (s *EtcdServer) requireAuthInfo(ctx context.Context) error { if !s.authStore.IsAuthEnabled() { return nil } authInfo, err := s.AuthInfoFromCtx(ctx) if err != nil { return err } if authInfo == nil { return auth.ErrUserEmpty } return nil }