From 9fe70138f00bc4cbfae647988775afe4ef0ece1f Mon Sep 17 00:00:00 2001 From: Andrey Butusov Date: Fri, 29 May 2026 15:31:04 +0300 Subject: [PATCH 1/2] node: inter-node mTLS Storage nodes authenticate to each other with mutual TLS over their existing public gRPC port. A node dials with a sentinel SNI and is served the identity certificate, verified against the network map; clients keep getting the plain or server-TLS endpoint unchanged. Peers are pinned by their network-map key. Signed-off-by: Andrey Butusov --- cmd/neofs-node/config.go | 16 +++- cmd/neofs-node/grpc.go | 28 ++---- cmd/neofs-node/mtls.go | 159 +++++++++++++++++++++++++++++++ pkg/network/cache/clients.go | 56 ++++++++--- pkg/network/peerauth/peerauth.go | 67 +++++++++++++ 5 files changed, 290 insertions(+), 36 deletions(-) create mode 100644 cmd/neofs-node/mtls.go create mode 100644 pkg/network/peerauth/peerauth.go diff --git a/cmd/neofs-node/config.go b/cmd/neofs-node/config.go index cba66be937..a979c21255 100644 --- a/cmd/neofs-node/config.go +++ b/cmd/neofs-node/config.go @@ -2,6 +2,8 @@ package main import ( "context" + "crypto/tls" + "encoding/hex" "errors" "fmt" "io/fs" @@ -33,6 +35,7 @@ import ( "github.com/nspcc-dev/neofs-node/pkg/morph/event" "github.com/nspcc-dev/neofs-node/pkg/network" "github.com/nspcc-dev/neofs-node/pkg/network/cache" + "github.com/nspcc-dev/neofs-node/pkg/network/peerauth" "github.com/nspcc-dev/neofs-node/pkg/services/control" controlSvc "github.com/nspcc-dev/neofs-node/pkg/services/control/server" "github.com/nspcc-dev/neofs-node/pkg/services/meta" @@ -160,6 +163,10 @@ type shared struct { putClientCache *cache.Clients localAddr network.AddressGroup + // tlsCert is the node's self-signed certificate built from its identity + // private key for inter-node mTLS handshakes. + tlsCert tls.Certificate + ownerIDFromKey user.ID // user ID calculated from key // current network map @@ -414,13 +421,19 @@ func initCfg(appCfg *config.Config) *cfg { fatalOnErr(err) basicSharedConfig := initBasics(c, key, persistate) + + tlsCert, err := peerauth.GenerateSelfSignedCert(&key.PrivateKey) + fatalOnErr(err) + c.log.Info("mTLS: generated self-signed certificate from node identity key", + zap.String("pubkey", hex.EncodeToString(key.PublicKey().Bytes()))) + streamTimeout := appCfg.APIClient.StreamTimeout minConnTimeout := appCfg.APIClient.MinConnectionTime pingInterval := appCfg.APIClient.PingInterval pingTimeout := appCfg.APIClient.PingTimeout newClientCache := func(scope string) *cache.Clients { return cache.NewClients(c.log.With(zap.String("scope", scope)), &buffers, streamTimeout, - minConnTimeout, pingInterval, pingTimeout, neofsecdsa.Signer(key.PrivateKey)) + minConnTimeout, pingInterval, pingTimeout, neofsecdsa.Signer(key.PrivateKey), tlsCert) } c.shared = shared{ basics: basicSharedConfig, @@ -430,6 +443,7 @@ func initCfg(appCfg *config.Config) *cfg { putClientCache: newClientCache("put"), persistate: persistate, privateTokenStore: persistate, + tlsCert: tlsCert, } c.cfgBalance = cfgBalance{ parsers: make(map[event.Type]event.NotificationParser), diff --git a/cmd/neofs-node/grpc.go b/cmd/neofs-node/grpc.go index 20008e933c..d7fa0fb958 100644 --- a/cmd/neofs-node/grpc.go +++ b/cmd/neofs-node/grpc.go @@ -19,7 +19,6 @@ import ( "go.uber.org/zap" "golang.org/x/net/netutil" "google.golang.org/grpc" - "google.golang.org/grpc/credentials" "google.golang.org/grpc/keepalive" "google.golang.org/grpc/resolver" ) @@ -185,30 +184,17 @@ func buildSingleGRPCServer(c *cfg, sc grpcconfig.GRPC, maxRecvMsgSizeOpt grpc.Se serverOpts = append(serverOpts, maxRecvMsgSizeOpt) } - tlsCfg := sc.TLS - - if tlsCfg.Key != "" { - certFile, keyFile := tlsCfg.Certificate, tlsCfg.Key - - if _, err := tls.LoadX509KeyPair(certFile, keyFile); err != nil { + if sc.TLS.Key != "" { + if _, err := tls.LoadX509KeyPair(sc.TLS.Certificate, sc.TLS.Key); err != nil { c.log.Error("could not read certificate from file", zap.Error(err)) return nil, nil, err } - // read certificate from disk on each handshake to pick up renewals automatically. - creds := credentials.NewTLS(&tls.Config{ - GetConfigForClient: func(*tls.ClientHelloInfo) (*tls.Config, error) { - cert, err := tls.LoadX509KeyPair(certFile, keyFile) - if err != nil { - return nil, fmt.Errorf("reload TLS certificate: %w", err) - } - return &tls.Config{ - Certificates: []tls.Certificate{cert}, - }, nil - }, - }) - - serverOpts = append(serverOpts, grpc.Creds(creds)) + serverOpts = append(serverOpts, grpc.Creds(newNodeAuthCreds(c, sc.TLS.Certificate, sc.TLS.Key))) + c.log.Info("gRPC endpoint serves clients (server-side TLS) and inter-node mTLS", zap.String("endpoint", sc.Endpoint)) + } else { + serverOpts = append(serverOpts, grpc.Creds(newNodeAuthCreds(c, "", ""))) + c.log.Info("gRPC endpoint serves clients (plain) and inter-node mTLS", zap.String("endpoint", sc.Endpoint)) } lis, err := net.Listen("tcp", sc.Endpoint) diff --git a/cmd/neofs-node/mtls.go b/cmd/neofs-node/mtls.go new file mode 100644 index 0000000000..29e79afcd1 --- /dev/null +++ b/cmd/neofs-node/mtls.go @@ -0,0 +1,159 @@ +package main + +import ( + "bytes" + "context" + "crypto/tls" + "crypto/x509" + "encoding/hex" + "fmt" + "io" + "net" + + "github.com/nspcc-dev/neofs-node/pkg/network/peerauth" + "github.com/nspcc-dev/neofs-sdk-go/netmap" + "go.uber.org/zap" + "google.golang.org/grpc/credentials" + "google.golang.org/grpc/credentials/insecure" +) + +// tlsRecordTypeHandshake is the first byte (ContentType) of a TLS handshake +// record. +const tlsRecordTypeHandshake = 0x16 + +// nodeAuthCreds are gRPC server transport credentials that serve regular clients +// and inter-node mTLS connections on the same port. +type nodeAuthCreds struct { + tls credentials.TransportCredentials + plain credentials.TransportCredentials +} + +// newNodeAuthCreds builds the shared-port credentials. adminCertFile and +// adminKeyFile are the optional server-side TLS certificate served to regular +// clients; when empty, only plain clients and inter-node mTLS are expected. +func newNodeAuthCreds(c *cfg, adminCertFile, adminKeyFile string) credentials.TransportCredentials { + return &nodeAuthCreds{ + tls: credentials.NewTLS(serverTLSConfig(c, adminCertFile, adminKeyFile)), + plain: insecure.NewCredentials(), + } +} + +func (a *nodeAuthCreds) ServerHandshake(raw net.Conn) (net.Conn, credentials.AuthInfo, error) { + first := make([]byte, 1) + if _, err := io.ReadFull(raw, first); err != nil { + return nil, nil, fmt.Errorf("peek first connection byte: %w", err) + } + conn := &prefixConn{Conn: raw, prefix: first} + if first[0] == tlsRecordTypeHandshake { + return a.tls.ServerHandshake(conn) + } + return a.plain.ServerHandshake(conn) +} + +func (a *nodeAuthCreds) ClientHandshake(ctx context.Context, authority string, raw net.Conn) (net.Conn, credentials.AuthInfo, error) { + return a.plain.ClientHandshake(ctx, authority, raw) +} + +func (a *nodeAuthCreds) Info() credentials.ProtocolInfo { return a.tls.Info() } + +func (a *nodeAuthCreds) Clone() credentials.TransportCredentials { + return &nodeAuthCreds{tls: a.tls.Clone(), plain: a.plain.Clone()} +} + +func (a *nodeAuthCreds) OverrideServerName(string) error { return nil } + +// prefixConn is a net.Conn that replays a number of already-read bytes before +// continuing with the underlying connection. It lets the handshake sniffer put +// the peeked byte back so the chosen handshaker sees the full stream. +type prefixConn struct { + net.Conn + prefix []byte +} + +func (c *prefixConn) Read(p []byte) (int, error) { + if len(c.prefix) > 0 { + n := copy(p, c.prefix) + c.prefix = c.prefix[n:] + return n, nil + } + return c.Conn.Read(p) +} + +func serverTLSConfig(c *cfg, adminCertFile, adminKeyFile string) *tls.Config { + node := nodeServerTLSConfig(c) + return &tls.Config{ + MinVersion: tls.VersionTLS12, + NextProtos: []string{"h2"}, + GetConfigForClient: func(hello *tls.ClientHelloInfo) (*tls.Config, error) { + if hello.ServerName == peerauth.TLSServerName { + return node, nil + } + if adminKeyFile == "" { + return nil, fmt.Errorf("no server TLS certificate configured for client %q", hello.ServerName) + } + cert, err := tls.LoadX509KeyPair(adminCertFile, adminKeyFile) + if err != nil { + return nil, fmt.Errorf("reload TLS certificate: %w", err) + } + return &tls.Config{ + Certificates: []tls.Certificate{cert}, + MinVersion: tls.VersionTLS12, + NextProtos: []string{"h2"}, + }, nil + }, + } +} + +// nodeServerTLSConfig builds the TLS sub-configuration used for inter-node mTLS +// connections. Only fellow storage nodes are expected here: the verifier +// rejects any peer whose key is not in the current network map. +func nodeServerTLSConfig(c *cfg) *tls.Config { + return &tls.Config{ + Certificates: []tls.Certificate{c.tlsCert}, + ClientAuth: tls.RequireAnyClientCert, + VerifyPeerCertificate: makeNetmapVerifier(c), + MinVersion: tls.VersionTLS12, + NextProtos: []string{"h2"}, + } +} + +// makeNetmapVerifier returns a TLS VerifyPeerCertificate hook that accepts a +// client certificate only if its public key belongs to a node in the current +// network map. +func makeNetmapVerifier(c *cfg) func(rawCerts [][]byte, _ [][]*x509.Certificate) error { + return func(rawCerts [][]byte, _ [][]*x509.Certificate) error { + pub, err := peerauth.PeerPubKey(rawCerts) + if err != nil { + return fmt.Errorf("extract peer pubkey: %w", err) + } + + if !isNetmapNode(c, pub) { + c.log.Warn("mTLS: rejecting inter-node peer absent from network map", + zap.String("pubkey", hex.EncodeToString(pub))) + return fmt.Errorf("peer key %s is not a network map node", hex.EncodeToString(pub)) + } + + c.log.Info("mTLS: peer authenticated as known SN", + zap.String("pubkey", hex.EncodeToString(pub))) + return nil + } +} + +// isNetmapNode reports whether pub is the public key of some node in the +// current network map snapshot. +func isNetmapNode(c *cfg, pub []byte) bool { + val := c.netMap.Load() + if val == nil { + return false + } + nm, ok := val.(netmap.NetMap) + if !ok { + return false + } + for _, n := range nm.Nodes() { + if bytes.Equal(n.PublicKey(), pub) { + return true + } + } + return false +} diff --git a/pkg/network/cache/clients.go b/pkg/network/cache/clients.go index 8da80f8909..d9d93ceae6 100644 --- a/pkg/network/cache/clients.go +++ b/pkg/network/cache/clients.go @@ -3,11 +3,12 @@ package cache import ( "bytes" "context" + "crypto/tls" + "crypto/x509" "encoding/hex" "errors" "fmt" "io" - "iter" "maps" "slices" "sync" @@ -17,6 +18,7 @@ import ( "github.com/nspcc-dev/neofs-node/internal/uriutil" clientcore "github.com/nspcc-dev/neofs-node/pkg/core/client" "github.com/nspcc-dev/neofs-node/pkg/network" + "github.com/nspcc-dev/neofs-node/pkg/network/peerauth" "github.com/nspcc-dev/neofs-sdk-go/client" cid "github.com/nspcc-dev/neofs-sdk-go/container/id" neofscrypto "github.com/nspcc-dev/neofs-sdk-go/crypto" @@ -29,7 +31,6 @@ import ( "google.golang.org/grpc" "google.golang.org/grpc/backoff" "google.golang.org/grpc/credentials" - "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/keepalive" ) @@ -48,12 +49,13 @@ type Clients struct { mtx sync.RWMutex conns map[string]*connections // keys are public key bytes - signer neofscrypto.Signer + signer neofscrypto.Signer + tlsCert tls.Certificate } // NewClients constructs Clients initializing connection to any endpoint with // given parameters. -func NewClients(l *zap.Logger, signBufPool *sync.Pool, streamTimeout, minConnTimeout, pingInterval, pingTimeout time.Duration, signer neofscrypto.Signer) *Clients { +func NewClients(l *zap.Logger, signBufPool *sync.Pool, streamTimeout, minConnTimeout, pingInterval, pingTimeout time.Duration, signer neofscrypto.Signer, tlsCert tls.Certificate) *Clients { return &Clients{ log: l, streamMsgTimeout: streamTimeout, @@ -63,6 +65,7 @@ func NewClients(l *zap.Logger, signBufPool *sync.Pool, streamTimeout, minConnTim pingTimeout: pingTimeout, conns: make(map[string]*connections), signer: signer, + tlsCert: tlsCert, } } @@ -96,7 +99,7 @@ func (x *Clients) Get(ctx context.Context, info netmap.NodeInfo) (clientcore.Mul return c, nil } - c, err := x.initConnections(ctx, info.PublicKey(), info.NetworkEndpoints()) + c, err := x.initConnections(ctx, info.PublicKey(), slices.Collect(info.NetworkEndpoints())) if err != nil { return nil, fmt.Errorf("init connections: %w", err) } @@ -109,10 +112,14 @@ func (x *Clients) SyncWithNewNetmap(ctx context.Context, sns []netmap.NodeInfo, x.mtx.Lock() defer x.mtx.Unlock() + localPub, _ := peerauth.CompressedPubKey(x.tlsCert.Leaf) for i := range sns { if i == local { continue } + if localPub != nil && bytes.Equal(sns[i].PublicKey(), localPub) { + continue + } if err := x.syncWithNetmapSN(ctx, sns[i]); err != nil { x.log.Warn("failed to sync connection cache with SN from the new network map, skip", zap.String("pub", hex.EncodeToString(sns[i].PublicKey())), zap.Error(err)) @@ -135,6 +142,11 @@ func (x *Clients) syncWithNetmapSN(ctx context.Context, sn netmap.NodeInfo) erro pub := sn.PublicKey() conns, ok := x.conns[snCacheKey(pub)] if !ok { + c, err := x.initConnections(ctx, pub, slices.Collect(sn.NetworkEndpoints())) + if err != nil { + return fmt.Errorf("eager init: %w", err) + } + x.conns[snCacheKey(pub)] = c return nil } @@ -175,10 +187,10 @@ func (x *Clients) syncWithNetmapSN(ctx context.Context, sn netmap.NodeInfo) erro return nil } -func (x *Clients) initConnections(ctx context.Context, pub []byte, addrs iter.Seq[string]) (*connections, error) { +func (x *Clients) initConnections(ctx context.Context, pub []byte, addrs []string) (*connections, error) { m := make(map[string]*client.Client) l := x.log.With(zap.String("public key", hex.EncodeToString(pub))) - for s := range addrs { + for _, s := range addrs { l.Info("initializing connection to the SN...", zap.String("address", s)) c, err := x.initConnection(ctx, pub, s) if err != nil { @@ -206,16 +218,32 @@ func (x *Clients) initConnection(ctx context.Context, pub []byte, uri string) (* return nil, fmt.Errorf("parse network address %q: %w", uri, err) } - target, withTLS, err := uriutil.Parse(a.URIAddr()) + target, _, err := uriutil.Parse(a.URIAddr()) if err != nil { return nil, fmt.Errorf("parse URI: %w", err) } - var transportCreds credentials.TransportCredentials - if withTLS { - transportCreds = credentials.NewTLS(nil) - } else { - transportCreds = insecure.NewCredentials() - } + transportCreds := credentials.NewTLS(&tls.Config{ + Certificates: []tls.Certificate{x.tlsCert}, + ServerName: peerauth.TLSServerName, // SNI marks this as an inter-node connection + InsecureSkipVerify: true, // CA chain is irrelevant; we pin by pubkey below + MinVersion: tls.VersionTLS12, + VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error { + got, err := peerauth.PeerPubKey(rawCerts) + if err != nil { + return fmt.Errorf("extract peer pubkey: %w", err) + } + if !bytes.Equal(got, pub) { + x.log.Warn("mTLS: server pubkey mismatch", + zap.String("expected", hex.EncodeToString(pub)), + zap.String("got", hex.EncodeToString(got))) + return clientcore.ErrWrongPublicKey + } + x.log.Info("mTLS: server verified by pubkey", + zap.String("pubkey", hex.EncodeToString(got)), + zap.String("target", target)) + return nil + }, + }) grpcConn, err := grpc.NewClient(target, grpc.WithTransportCredentials(transportCreds), grpc.WithConnectParams(grpc.ConnectParams{ diff --git a/pkg/network/peerauth/peerauth.go b/pkg/network/peerauth/peerauth.go new file mode 100644 index 0000000000..b8bdec9426 --- /dev/null +++ b/pkg/network/peerauth/peerauth.go @@ -0,0 +1,67 @@ +package peerauth + +import ( + "crypto/ecdsa" + "crypto/rand" + "crypto/tls" + "crypto/x509" + "errors" + "fmt" + "math/big" + "time" + + "github.com/nspcc-dev/neo-go/pkg/crypto/keys" +) + +// TLSServerName is the TLS SNI a storage node sends when dialing a peer over the +// shared public port. The server uses it to distinguish inter-node mTLS +// connections from regular client TLS connections. +const TLSServerName = "neofs.internode.mtls" + +// GenerateSelfSignedCert builds a self-signed TLS certificate that uses the +// node's identity ECDSA key as both subject key and signing key. +func GenerateSelfSignedCert(priv *ecdsa.PrivateKey) (tls.Certificate, error) { + template := &x509.Certificate{ + SerialNumber: big.NewInt(1), + NotBefore: time.Now().Add(-time.Minute), + NotAfter: time.Now().Add(100 * 365 * 24 * time.Hour), + KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageKeyEncipherment, + ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth, x509.ExtKeyUsageClientAuth}, + } + der, err := x509.CreateCertificate(rand.Reader, template, template, &priv.PublicKey, priv) + if err != nil { + return tls.Certificate{}, fmt.Errorf("create x509 certificate: %w", err) + } + leaf, err := x509.ParseCertificate(der) + if err != nil { + return tls.Certificate{}, fmt.Errorf("parse generated certificate: %w", err) + } + return tls.Certificate{ + Certificate: [][]byte{der}, + PrivateKey: priv, + Leaf: leaf, + }, nil +} + +// CompressedPubKey returns the compressed-form public key bytes of an ECDSA +// certificate, matching the format used in NeoFS network maps. +func CompressedPubKey(cert *x509.Certificate) ([]byte, error) { + pub, ok := cert.PublicKey.(*ecdsa.PublicKey) + if !ok { + return nil, errors.New("peer certificate public key is not ECDSA") + } + return (*keys.PublicKey)(pub).Bytes(), nil +} + +// PeerPubKey parses the first certificate from a TLS handshake's raw chain and +// returns the peer's compressed public key bytes. +func PeerPubKey(rawCerts [][]byte) ([]byte, error) { + if len(rawCerts) == 0 { + return nil, errors.New("no peer certificate") + } + cert, err := x509.ParseCertificate(rawCerts[0]) + if err != nil { + return nil, fmt.Errorf("parse peer certificate: %w", err) + } + return CompressedPubKey(cert) +} From 810ccfed4caf2baa2dcfbe8831c65edee6ecdf71 Mon Sep 17 00:00:00 2001 From: Andrey Butusov Date: Fri, 29 May 2026 15:32:27 +0300 Subject: [PATCH 2/2] object: skip request/response signing on authenticated 1:1 node hops On the inter-node mTLS listener the peer is an authenticated network-map node, so verifying the per-request signature and signing the response are redundant for 1:1 (TTL<=1) hops and are skipped. Requests arriving on the plain public listener carry no TLS peer certificate and are always verified and signed, so clients remain unaffected. Signed-off-by: Andrey Butusov --- pkg/services/object/acl/v2/service.go | 14 +++ pkg/services/object/common/request.go | 4 + pkg/services/object/server.go | 170 +++++++++++++++++--------- 3 files changed, 133 insertions(+), 55 deletions(-) diff --git a/pkg/services/object/acl/v2/service.go b/pkg/services/object/acl/v2/service.go index 4c177bc45e..1ea452e4fc 100644 --- a/pkg/services/object/acl/v2/service.go +++ b/pkg/services/object/acl/v2/service.go @@ -1,6 +1,8 @@ package v2 import ( + "crypto/ecdsa" + "crypto/elliptic" "crypto/sha256" "errors" "fmt" @@ -9,6 +11,7 @@ import ( lru "github.com/hashicorp/golang-lru/v2" "github.com/nspcc-dev/neo-go/pkg/core/block" "github.com/nspcc-dev/neo-go/pkg/core/transaction" + "github.com/nspcc-dev/neo-go/pkg/crypto/keys" "github.com/nspcc-dev/neo-go/pkg/neorpc/result" "github.com/nspcc-dev/neo-go/pkg/smartcontract/trigger" "github.com/nspcc-dev/neo-go/pkg/util" @@ -293,6 +296,15 @@ func getCredentialsFromSessionToken(token sessionv2.Token) (user.ID, []byte, err return token.OriginalIssuer(), key, nil } +func getCredentialsFromPeerPublicKey(key []byte) (user.ID, []byte, error) { + pub, err := keys.NewPublicKeyFromBytes(key, elliptic.P256()) + if err != nil { + return user.ID{}, nil, fmt.Errorf("invalid peer public key: %w", err) + } + + return user.NewFromECDSAPublicKey(ecdsa.PublicKey(*pub)), key, nil +} + type sessionTokenV2WithEncodedBody struct { sessionv2.Token body []byte @@ -505,6 +517,8 @@ func (b Service) findRequestInfo(req interface { reqAuthor, reqAuthorPub, err = getCredentialsFromSessionToken(*tokens.Session) } else if tokens.SessionV1 != nil { reqAuthor, reqAuthorPub, err = getCredentialsFromSessionV1Token(*tokens.SessionV1) + } else if tokens.AuthenticatedPeerPublicKey != nil { + reqAuthor, reqAuthorPub, err = getCredentialsFromPeerPublicKey(tokens.AuthenticatedPeerPublicKey) } else { reqAuthor, reqAuthorPub, err = icrypto.GetRequestAuthor(req.GetVerifyHeader()) } diff --git a/pkg/services/object/common/request.go b/pkg/services/object/common/request.go index f3dc53671b..d60af96230 100644 --- a/pkg/services/object/common/request.go +++ b/pkg/services/object/common/request.go @@ -11,4 +11,8 @@ type RequestTokens struct { Session *sessionv2.Token SessionV1 *session.Object Bearer *bearer.Token + + // AuthenticatedPeerPublicKey is the compressed ECDSA public key of an + // inter-node mTLS peer authenticated by the transport layer. + AuthenticatedPeerPublicKey []byte } diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index 586e597441..cb1c36c6e1 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -26,6 +26,7 @@ import ( containercore "github.com/nspcc-dev/neofs-node/pkg/core/container" netmapcore "github.com/nspcc-dev/neofs-node/pkg/core/netmap" objectcore "github.com/nspcc-dev/neofs-node/pkg/core/object" + "github.com/nspcc-dev/neofs-node/pkg/network/peerauth" metasvc "github.com/nspcc-dev/neofs-node/pkg/services/meta" aclsvc "github.com/nspcc-dev/neofs-node/pkg/services/object/acl/v2" "github.com/nspcc-dev/neofs-node/pkg/services/object/common" @@ -59,7 +60,9 @@ import ( "go.uber.org/zap" "google.golang.org/grpc" grpccodes "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials" "google.golang.org/grpc/mem" + "google.golang.org/grpc/peer" grpcstatus "google.golang.org/grpc/status" "google.golang.org/protobuf/proto" ) @@ -265,7 +268,9 @@ func (s *Server) sendPutResponse(stream protoobject.ObjectService_PutServer, res resp.MetaHeader = s.makeResponseMetaHeader(util.ToStatus(err)) } - resp.VerifyHeader = util.SignResponseIfNeeded(&s.signer, resp, req) + if needSignObjectResponse(stream.Context(), req) { + resp.VerifyHeader = util.SignResponse(&s.signer, resp) + } return stream.SendAndClose(resp) } @@ -335,10 +340,12 @@ func (x *putStream) resignRequest(req *protoobject.PutRequest) (*protoobject.Put Ttl: meta.GetTtl() - 1, Origin: meta, } - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(x.signer), req, nil) - if err != nil { - return nil, err + if req.MetaHeader.Ttl > 1 { + var err error + req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(x.signer), req, nil) + if err != nil { + return nil, err + } } return req, nil } @@ -450,9 +457,11 @@ func (s *Server) Put(gStream protoobject.ObjectService_PutServer) error { s.metrics.AddPutPayload(len(c)) } - if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { - err = s.sendStatusPutResponse(gStream, err, reqFirst) // assign for defer - return err + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(gStream.Context()) { + if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { + err = s.sendStatusPutResponse(gStream, err, reqFirst) // assign for defer + return err + } } if s.fsChain.LocalNodeUnderMaintenance() { @@ -519,7 +528,7 @@ func (s *Server) Put(gStream protoobject.ObjectService_PutServer) error { return s.sendStatusPutResponse(gStream, err, reqFirst) } - if reqInfo, objOwner, err := s.reqInfoProc.PutRequestToInfo(req, initPart, cnrID, op, reqMD.tokens); err != nil { + if reqInfo, objOwner, err := s.reqInfoProc.PutRequestToInfo(req, initPart, cnrID, op, requestTokensWithPeer(gStream.Context(), req, reqMD.tokens)); err != nil { if !errors.Is(err, aclsvc.ErrSkipRequest) { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -545,16 +554,18 @@ func (s *Server) Put(gStream protoobject.ObjectService_PutServer) error { } } -func (s *Server) signDeleteResponse(resp *protoobject.DeleteResponse, err error, req *protoobject.DeleteRequest) *protoobject.DeleteResponse { +func (s *Server) signDeleteResponse(resp *protoobject.DeleteResponse, err error, sign bool) *protoobject.DeleteResponse { if err != nil { resp.MetaHeader = s.makeResponseMetaHeader(util.ToStatus(err)) } - resp.VerifyHeader = util.SignResponseIfNeeded(&s.signer, resp, req) + if sign { + resp.VerifyHeader = util.SignResponse(&s.signer, resp) + } return resp } -func (s *Server) makeStatusDeleteResponse(err error, req *protoobject.DeleteRequest) *protoobject.DeleteResponse { - return s.signDeleteResponse(new(protoobject.DeleteResponse), err, req) +func (s *Server) makeStatusDeleteResponse(err error, sign bool) *protoobject.DeleteResponse { + return s.signDeleteResponse(new(protoobject.DeleteResponse), err, sign) } type deleteResponseBody protoobject.DeleteResponse_Body @@ -570,46 +581,50 @@ func (s *Server) Delete(ctx context.Context, req *protoobject.DeleteRequest) (*p ) defer func() { s.pushOpExecResult(stat.MethodObjectDelete, err, t) }() - if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { - return s.makeStatusDeleteResponse(err, req), nil + needSignResp := needSignObjectResponse(ctx, req) + + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(ctx) { + if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { + return s.makeStatusDeleteResponse(err, needSignResp), nil + } } if s.fsChain.LocalNodeUnderMaintenance() { - return s.makeStatusDeleteResponse(apistatus.ErrNodeUnderMaintenance, req), nil + return s.makeStatusDeleteResponse(apistatus.ErrNodeUnderMaintenance, needSignResp), nil } body := req.Body if body == nil { err = newBadRequestError(missingRequestBodyMessage) // defer - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } cnrID, objID, err := fetchRequiredObjectAddress(body.Address) if err != nil { err = newBadRequestError(invalidRequestBodyMessage + ": " + err.Error()) // defer - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } reqMD, err := s.handleRequestMetaHeader(req.MetaHeader, sessionv2.VerbObjectDelete, session.VerbObjectDelete, cnrID, objID) if err != nil { - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } - reqInfo, err := s.reqInfoProc.DeleteRequestToInfo(req, cnrID, reqMD.tokens) + reqInfo, err := s.reqInfoProc.DeleteRequestToInfo(req, cnrID, requestTokensWithPeer(ctx, req, reqMD.tokens)) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer } - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } if !s.aclChecker.CheckBasicACL(reqInfo) { err = basicACLErr(reqInfo) // needed for defer - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } err = s.aclChecker.CheckEACL(ctx, req, cnrID, objID, reqInfo) if err != nil && !errors.Is(err, aclsvc.ErrNotMatched) { // Not matched -> follow basic ACL. err = eACLErr(reqInfo, err) // needed for defer - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } cp := objutil.CommonPrmFromRequest(reqMD.ttl, reqMD.xHeaders, reqMD.tokens) @@ -622,10 +637,10 @@ func (s *Server) Delete(ctx context.Context, req *protoobject.DeleteRequest) (*p p.WithTombstoneAddressTarget((*deleteResponseBody)(&rb)) err = s.handlers.Delete(ctx, p) if err != nil && !errors.Is(err, apistatus.ErrIncomplete) { - return s.makeStatusDeleteResponse(err, req), nil + return s.makeStatusDeleteResponse(err, needSignResp), nil } - return s.signDeleteResponse(&protoobject.DeleteResponse{Body: &rb}, err, req), nil + return s.signDeleteResponse(&protoobject.DeleteResponse{Body: &rb}, err, needSignResp), nil } func (s *Server) signHeadResponse(resp *protoobject.HeadResponse, sign bool) *protoobject.HeadResponse { @@ -669,10 +684,12 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) ) defer func() { s.pushOpExecResult(stat.MethodObjectHead, err, t) }() - needSignResp := needSignGetResponse(req) + needSignResp := needSignGetResponse(ctx, req) - if err := icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { - return s.makeStatusHeadResponse(err, needSignResp) + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(ctx) { + if err := icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { + return s.makeStatusHeadResponse(err, needSignResp) + } } if s.fsChain.LocalNodeUnderMaintenance() { @@ -696,7 +713,7 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) return s.makeStatusHeadResponse(err, needSignResp) } - reqInfo, err := s.reqInfoProc.HeadRequestToInfo(req, cnrID, reqMD.tokens) + reqInfo, err := s.reqInfoProc.HeadRequestToInfo(req, cnrID, requestTokensWithPeer(ctx, req, reqMD.tokens)) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -864,11 +881,6 @@ func convertHeadPrm(signer ecdsa.PrivateKey, cnr container.Container, req *proto Ttl: 1, }, } - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) - if err != nil { - return nil, iprotobuf.BuffersSlice{}, err - } } var respBuf mem.BufferSlice @@ -991,10 +1003,12 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ ) defer func() { s.pushOpExecResult(stat.MethodObjectGet, err, t) }() - needSignResp := needSignGetResponse(req) + needSignResp := needSignGetResponse(gStream.Context(), req) - if err = icrypto.VerifyRequestSignatures(req); err != nil { - return s.sendStatusGetResponse(gStream, err, needSignResp) + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(gStream.Context()) { + if err = icrypto.VerifyRequestSignatures(req); err != nil { + return s.sendStatusGetResponse(gStream, err, needSignResp) + } } if s.fsChain.LocalNodeUnderMaintenance() { @@ -1018,7 +1032,7 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ return s.sendStatusGetResponse(gStream, err, needSignResp) } - reqInfo, err := s.reqInfoProc.GetRequestToInfo(req, cnrID, reqMD.tokens) + reqInfo, err := s.reqInfoProc.GetRequestToInfo(req, cnrID, requestTokensWithPeer(gStream.Context(), req, reqMD.tokens)) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -1299,11 +1313,6 @@ func convertGetPrm(signer ecdsa.PrivateKey, cnr container.Container, req *protoo if proxyCtx.suppressInit { req.Body.PayloadOnly = false } - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) - if err != nil { - return err - } } return c.ForAnyGRPCConn(ctx, func(ctx context.Context, conn *grpc.ClientConn) error { @@ -1379,8 +1388,10 @@ func (s *Server) GetRange(req *protoobject.GetRangeRequest, gStream protoobject. t = time.Now() ) defer func() { s.pushOpExecResult(stat.MethodObjectRange, err, t) }() - if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { - return s.sendStatusRangeResponse(gStream, err, req) + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(gStream.Context()) { + if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { + return s.sendStatusRangeResponse(gStream, err, req) + } } if s.fsChain.LocalNodeUnderMaintenance() { @@ -1404,7 +1415,7 @@ func (s *Server) GetRange(req *protoobject.GetRangeRequest, gStream protoobject. return s.sendStatusRangeResponse(gStream, err, req) } - reqInfo, err := s.reqInfoProc.RangeRequestToInfo(req, cnrID, reqMD.tokens) + reqInfo, err := s.reqInfoProc.RangeRequestToInfo(req, cnrID, requestTokensWithPeer(gStream.Context(), req, reqMD.tokens)) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -1421,7 +1432,7 @@ func (s *Server) GetRange(req *protoobject.GetRangeRequest, gStream protoobject. return s.sendStatusRangeResponse(gStream, err, req) } - needSignResponse := needSignGetResponse(req) + needSignResponse := needSignGetResponse(gStream.Context(), req) p, err := convertRangePrm(s.signer, reqInfo.Container, req, &rangeStream{ base: gStream, @@ -1558,7 +1569,9 @@ func convertRangePrm(signer ecdsa.PrivateKey, cnr container.Container, req *prot Ttl: meta.GetTtl() - 1, Origin: meta, } - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) + if req.MetaHeader.Ttl > 1 { + req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) + } }) if err != nil { return err @@ -1775,8 +1788,10 @@ func (s *Server) SearchV2Buffered(ctx context.Context, req *protoobject.SearchV2 t = time.Now() ) defer s.pushOpExecResult(stat.MethodObjectSearchV2, err, t) - if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { - return s.signSearchResponse(nil, err, req) + if req.GetMetaHeader().GetTtl() > 1 || !peerAuthenticatedAsNode(ctx) { + if err = icrypto.VerifyRequestSignaturesN3(req, s.fsChain); err != nil { + return s.signSearchResponse(nil, err, req) + } } if s.fsChain.LocalNodeUnderMaintenance() { @@ -1800,7 +1815,7 @@ func (s *Server) SearchV2Buffered(ctx context.Context, req *protoobject.SearchV2 return s.signSearchResponse(nil, err, req) } - reqInfo, err := s.reqInfoProc.SearchV2RequestToInfo(req, cnrID, reqMD.tokens) + reqInfo, err := s.reqInfoProc.SearchV2RequestToInfo(req, cnrID, requestTokensWithPeer(ctx, req, reqMD.tokens)) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -1991,9 +2006,6 @@ func (s *Server) ProcessSearch(ctx context.Context, req *protoobject.SearchV2Req Ttl: 1, }, } - if req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer[*protoobject.SearchV2Request_Body](neofsecdsa.Signer(s.signer), req, nil); err != nil { - return nil, nil, fmt.Errorf("sign request: %w", err) - } var optimizedNodes = (len(body.Filters) != 0) && slices.ContainsFunc(body.Filters, func(filt *protoobject.SearchFilter) bool { @@ -2244,10 +2256,58 @@ func chunkBoundsToSend(global, local, chunkLen int) (int, int) { return global - local, chunkLen } -func needSignGetResponse(req util.Request) bool { +func needSignObjectResponse(ctx context.Context, req util.Request) bool { + if req.GetMetaHeader().GetTtl() <= 1 && peerAuthenticatedAsNode(ctx) { + return false + } + return util.VersionLE(req, 2, 21) +} + +func needSignGetResponse(ctx context.Context, req util.Request) bool { + if req.GetMetaHeader().GetTtl() <= 1 && peerAuthenticatedAsNode(ctx) { + return false + } return util.VersionLE(req, 2, 17) } +// peerAuthenticatedAsNode reports whether the request arrived over the inter-node +// mTLS listener: a TLS peer certificate means the handshake verifier already +// matched the peer against the network map, so it is a trusted storage node. +// Only such 1:1 (TTL<=1) hops may skip signature checks; plain-listener clients +// carry no peer certificate and are always verified. +func peerAuthenticatedAsNode(ctx context.Context) bool { + _, ok := authenticatedPeerPublicKey(ctx) + return ok +} + +func authenticatedPeerPublicKey(ctx context.Context) ([]byte, bool) { + p, ok := peer.FromContext(ctx) + if !ok { + return nil, false + } + tlsInfo, ok := p.AuthInfo.(credentials.TLSInfo) + if !ok { + return nil, false + } + if len(tlsInfo.State.PeerCertificates) == 0 { + return nil, false + } + pubKey, err := peerauth.CompressedPubKey(tlsInfo.State.PeerCertificates[0]) + if err != nil { + return nil, false + } + return pubKey, true +} + +func requestTokensWithPeer(ctx context.Context, req util.Request, tokens common.RequestTokens) common.RequestTokens { + if meta := req.GetMetaHeader(); meta != nil && meta.GetTtl() <= 1 { + if pubKey, ok := authenticatedPeerPublicKey(ctx); ok { + tokens.AuthenticatedPeerPublicKey = pubKey + } + } + return tokens +} + func checkHeaderProtobufAgainstID(buffers iprotobuf.BuffersSlice, id oid.ID, ordered bool) error { b := buffers.ReadOnlyData() if !ordered {