From 41ad17b4b35479b8da551d6ba059958dacb8934e Mon Sep 17 00:00:00 2001 From: Leonard Liubich Date: Mon, 31 Aug 2026 18:16:08 +0300 Subject: [PATCH 1/3] sn/object: Deduplicate request encoding on GET/HEAD forwarding Follow 60d9d35be282f97b497d056875e7e6ba76aaf47b. Will be useful for #4078. Signed-off-by: Leonard Liubich --- CHANGELOG.md | 2 +- pkg/services/object/head.go | 2 ++ pkg/services/object/proto.go | 18 ++++++++++++++++++ pkg/services/object/search.go | 13 +------------ pkg/services/object/server.go | 26 ++++++++++++++++++++++++-- 5 files changed, 46 insertions(+), 15 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 48649313c0..4f435e888a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,7 +26,7 @@ Changelog for NeoFS Node - SNs no longer sign TTL=1 requests over mutually authenticated inter-node connections (#4100, #4159) - Optimized GRPC read/write bufferring (#4130) - TLS key of SN is now always read from the node wallet instead of configuration (#4131) -- SN now allocates less to forward SEARCH requests (#4116) +- SN now allocates less to forward SEARCH/GET/HEAD requests (#4116, #4160) - SN object GET server now communicates over remote SN API version if it is lower than the local one (#4153) ### Removed diff --git a/pkg/services/object/head.go b/pkg/services/object/head.go index ec07551611..f93f4a1dc3 100644 --- a/pkg/services/object/head.go +++ b/pkg/services/object/head.go @@ -19,6 +19,8 @@ import ( "google.golang.org/protobuf/encoding/protowire" ) +var headRequestBufferPool = mem.DefaultBufferPool() + func callHead(ctx context.Context, conn *grpc.ClientConn, req any) (mem.BufferSlice, error) { return callUnary(ctx, conn, protoobject.ObjectService_Head_FullMethodName, req) } diff --git a/pkg/services/object/proto.go b/pkg/services/object/proto.go index f5c354a129..3998f6d6b4 100644 --- a/pkg/services/object/proto.go +++ b/pkg/services/object/proto.go @@ -15,6 +15,7 @@ import ( iprotobuf "github.com/nspcc-dev/neofs-sdk-go/proto/protobuf" protosession "github.com/nspcc-dev/neofs-sdk-go/proto/session" "github.com/nspcc-dev/neofs-sdk-go/version" + "google.golang.org/grpc/mem" "google.golang.org/protobuf/encoding/protowire" ) @@ -340,3 +341,20 @@ func (s *Server) writeRequestSignatures(reqBuf []byte, bodyWithMetaLen int, body return nil } + +func encodeRequestProtobuf(pool mem.BufferPool, body protoencoding.Message, metaHdr protoencoding.Message, verifHdr protoencoding.Message) *[]byte { + bodyLen := body.MarshaledSize() + metaHdrLen := metaHdr.MarshaledSize() + verifHdrLen := verifHdr.MarshaledSize() + + reqLen := protoencoding.CalculateRequestLength(bodyLen, metaHdrLen, verifHdrLen) + + bufItem := pool.Get(reqLen) + + buf := *bufItem + off := protoencoding.WriteRequestBodyMessage(buf, body) + off += protoencoding.WriteRequestMetaHeaderMessage(buf[off:], metaHdr) + protoencoding.WriteRequestVerificationHeaderMessage(buf[off:], verifHdr) + + return bufItem +} diff --git a/pkg/services/object/search.go b/pkg/services/object/search.go index 070fe2c5b7..1e1542d3ee 100644 --- a/pkg/services/object/search.go +++ b/pkg/services/object/search.go @@ -12,7 +12,6 @@ import ( islices "github.com/nspcc-dev/neofs-node/internal/slices" clientcore "github.com/nspcc-dev/neofs-node/pkg/core/client" "github.com/nspcc-dev/neofs-sdk-go/netmap" - protoencoding "github.com/nspcc-dev/neofs-sdk-go/proto/encoding" protoobject "github.com/nspcc-dev/neofs-sdk-go/proto/object" "github.com/nspcc-dev/neofs-sdk-go/version" "go.uber.org/zap" @@ -53,20 +52,10 @@ func iterateSearchableContainerNodes(nodeSets [][]netmap.NodeInfo, repRules []ui } func (s *Server) forwardSearchRequest(ctx context.Context, req *protoobject.SearchV2Request, nodeSets [][]netmap.NodeInfo) (mem.BufferSlice, error) { - bodyLen := req.Body.MarshaledSize() - metaHdrLen := req.MetaHeader.MarshaledSize() - verifHdrLen := req.VerifyHeader.MarshaledSize() - - reqLen := protoencoding.CalculateRequestLength(bodyLen, metaHdrLen, verifHdrLen) - - bufItem := searchRequestBufferPool.Get(reqLen) + bufItem := encodeRequestProtobuf(searchRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) defer searchRequestBufferPool.Put(bufItem) buf := *bufItem - off := protoencoding.WriteRequestBodyMessage(buf, req.Body) - off += protoencoding.WriteRequestMetaHeaderMessage(buf[off:], req.MetaHeader) - protoencoding.WriteRequestVerificationHeaderMessage(buf[off:], req.VerifyHeader) - reqBuf := mem.SliceBuffer(buf) for _, nodeSet := range nodeSets { diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index 143212fbe5..6ae2b4b02c 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -713,9 +713,20 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) var forwardResp mem.BufferSlice + var forwardReqBufItem *[]byte + defer func() { + if forwardReqBufItem != nil { + headRequestBufferPool.Put(forwardReqBufItem) + } + }() + p.SetForwardRequestFunc(func(ctx context.Context, node clientcore.MultiAddressClient) error { + if forwardReqBufItem == nil { + forwardReqBufItem = encodeRequestProtobuf(headRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) + } + var err error - forwardResp, err = forwardHeadRequest(ctx, req, node) + forwardResp, err = forwardHeadRequest(ctx, mem.SliceBuffer(*forwardReqBufItem), node) return err }) @@ -1083,8 +1094,19 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ return s.sendStatusGetResponse(req, gStream, err, needSignResp) } + var forwardReqBufItem *[]byte + defer func() { + if forwardReqBufItem != nil { + getRequestBufferPool.Put(forwardReqBufItem) + } + }() + p.SetForwardRequestFunc(func(ctx context.Context, node clientcore.MultiAddressClient) error { - return forwardGetRequest(ctx, req, gStream, node) + if forwardReqBufItem == nil { + forwardReqBufItem = encodeRequestProtobuf(getRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) + } + + return forwardGetRequest(ctx, mem.SliceBuffer(*forwardReqBufItem), gStream, node) }) p.WithECTransport(&getECTransport{ From 743ee89e88f69ea6c1e3c9e5156102db64ba550e Mon Sep 17 00:00:00 2001 From: Leonard Liubich Date: Tue, 1 Sep 2026 09:30:15 +0300 Subject: [PATCH 2/3] sn/object: Deduplicate vars referring the same buffer pool So that it would be clear from code that the same pool is being used. Signed-off-by: Leonard Liubich --- pkg/services/object/get.go | 10 ++++------ pkg/services/object/head.go | 2 -- pkg/services/object/proto.go | 6 ++++-- pkg/services/object/search.go | 6 ++---- pkg/services/object/server.go | 12 ++++++------ 5 files changed, 16 insertions(+), 20 deletions(-) diff --git a/pkg/services/object/get.go b/pkg/services/object/get.go index 0c44857f11..a466cdd25b 100644 --- a/pkg/services/object/get.go +++ b/pkg/services/object/get.go @@ -33,8 +33,6 @@ import ( "google.golang.org/protobuf/encoding/protowire" ) -var getRequestBufferPool = mem.DefaultBufferPool() - var ( getStreamDesc = &grpc.StreamDesc{ StreamName: "Get", @@ -542,12 +540,12 @@ func (x *getECTransport) CopyECParentHeaderAndPayloadFromRemoteFirstPart(ctx con copiedHdr, parentPldLen, partPldLen, copiedPartPld, err = x.copyRemotePart(ctx, conn, req) if err != nil { - getRequestBufferPool.Put(x.getPartRequest) + defaultGRPCBufferPool.Put(x.getPartRequest) return err } if copiedHdr { - getRequestBufferPool.Put(x.getPartRequest) + defaultGRPCBufferPool.Put(x.getPartRequest) if copiedPartPld == partPldLen { return nil } @@ -696,7 +694,7 @@ func (x *getECTransport) copyRemotePartRange(ctx context.Context, conn *grpc.Cli copied, err := x.copyRemotePartRangeWithRequest(ctx, conn, request, ln, controlCh) if err != nil || copied == ln { - getRequestBufferPool.Put(request) + defaultGRPCBufferPool.Put(request) } return copied, err @@ -934,7 +932,7 @@ func (s *Server) makeGetECPartRequest(needSign bool, remoteServerAPIVersion vers reqLen += calculateRequestVerificationHeaderFieldLen(remoteServerAPIVersion) } - bufItem := getRequestBufferPool.Get(reqLen) + bufItem := defaultGRPCBufferPool.Get(reqLen) buf := *bufItem // body diff --git a/pkg/services/object/head.go b/pkg/services/object/head.go index f93f4a1dc3..ec07551611 100644 --- a/pkg/services/object/head.go +++ b/pkg/services/object/head.go @@ -19,8 +19,6 @@ import ( "google.golang.org/protobuf/encoding/protowire" ) -var headRequestBufferPool = mem.DefaultBufferPool() - func callHead(ctx context.Context, conn *grpc.ClientConn, req any) (mem.BufferSlice, error) { return callUnary(ctx, conn, protoobject.ObjectService_Head_FullMethodName, req) } diff --git a/pkg/services/object/proto.go b/pkg/services/object/proto.go index 3998f6d6b4..3368768b12 100644 --- a/pkg/services/object/proto.go +++ b/pkg/services/object/proto.go @@ -43,6 +43,8 @@ const ( responseVerificationHeaderECDSAWIthSHA512Len = verificationHeaderECDSAWithSHA512SignatureLen * 3 ) +var defaultGRPCBufferPool = mem.DefaultBufferPool() + var currentVersionResponseMetaHeader []byte func init() { @@ -342,14 +344,14 @@ func (s *Server) writeRequestSignatures(reqBuf []byte, bodyWithMetaLen int, body return nil } -func encodeRequestProtobuf(pool mem.BufferPool, body protoencoding.Message, metaHdr protoencoding.Message, verifHdr protoencoding.Message) *[]byte { +func encodeRequestProtobuf(body protoencoding.Message, metaHdr protoencoding.Message, verifHdr protoencoding.Message) *[]byte { bodyLen := body.MarshaledSize() metaHdrLen := metaHdr.MarshaledSize() verifHdrLen := verifHdr.MarshaledSize() reqLen := protoencoding.CalculateRequestLength(bodyLen, metaHdrLen, verifHdrLen) - bufItem := pool.Get(reqLen) + bufItem := defaultGRPCBufferPool.Get(reqLen) buf := *bufItem off := protoencoding.WriteRequestBodyMessage(buf, body) diff --git a/pkg/services/object/search.go b/pkg/services/object/search.go index 1e1542d3ee..ef4580c3c3 100644 --- a/pkg/services/object/search.go +++ b/pkg/services/object/search.go @@ -19,8 +19,6 @@ import ( "google.golang.org/grpc/mem" ) -var searchRequestBufferPool = mem.DefaultBufferPool() - func iterateSearchableContainerNodes(nodeSets [][]netmap.NodeInfo, repRules []uint, ecRules []iec.Rule, allNodes bool, f func(netmap.NodeInfo) bool) { for i := range nodeSets { var ( @@ -52,8 +50,8 @@ func iterateSearchableContainerNodes(nodeSets [][]netmap.NodeInfo, repRules []ui } func (s *Server) forwardSearchRequest(ctx context.Context, req *protoobject.SearchV2Request, nodeSets [][]netmap.NodeInfo) (mem.BufferSlice, error) { - bufItem := encodeRequestProtobuf(searchRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) - defer searchRequestBufferPool.Put(bufItem) + bufItem := encodeRequestProtobuf(req.Body, req.MetaHeader, req.VerifyHeader) + defer defaultGRPCBufferPool.Put(bufItem) buf := *bufItem reqBuf := mem.SliceBuffer(buf) diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index 6ae2b4b02c..7476cac9db 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -716,13 +716,13 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) var forwardReqBufItem *[]byte defer func() { if forwardReqBufItem != nil { - headRequestBufferPool.Put(forwardReqBufItem) + defaultGRPCBufferPool.Put(forwardReqBufItem) } }() p.SetForwardRequestFunc(func(ctx context.Context, node clientcore.MultiAddressClient) error { if forwardReqBufItem == nil { - forwardReqBufItem = encodeRequestProtobuf(headRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) + forwardReqBufItem = encodeRequestProtobuf(req.Body, req.MetaHeader, req.VerifyHeader) } var err error @@ -1097,13 +1097,13 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ var forwardReqBufItem *[]byte defer func() { if forwardReqBufItem != nil { - getRequestBufferPool.Put(forwardReqBufItem) + defaultGRPCBufferPool.Put(forwardReqBufItem) } }() p.SetForwardRequestFunc(func(ctx context.Context, node clientcore.MultiAddressClient) error { if forwardReqBufItem == nil { - forwardReqBufItem = encodeRequestProtobuf(getRequestBufferPool, req.Body, req.MetaHeader, req.VerifyHeader) + forwardReqBufItem = encodeRequestProtobuf(req.Body, req.MetaHeader, req.VerifyHeader) } return forwardGetRequest(ctx, mem.SliceBuffer(*forwardReqBufItem), gStream, node) @@ -2171,8 +2171,8 @@ func (s *Server) ProcessSearch(ctx context.Context, req *protoobject.SearchV2Req verifHdrFLdLen := calculateRequestVerificationHeaderFieldLen(localVersion) reqLen := bodyWithMetaLen + verifHdrFLdLen - reqBufItem := searchRequestBufferPool.Get(reqLen) - defer searchRequestBufferPool.Put(reqBufItem) + reqBufItem := defaultGRPCBufferPool.Get(reqLen) + defer defaultGRPCBufferPool.Put(reqBufItem) reqBuf := *reqBufItem off := protoencoding.WriteRequestBodyTagAndLength(reqBuf, bodyLen) From 8917f5b4715316a805fd9ed059538afb8d0e1927 Mon Sep 17 00:00:00 2001 From: Leonard Liubich Date: Tue, 1 Sep 2026 11:02:11 +0300 Subject: [PATCH 3/3] sn/object: Deduplicate encoding of SN-SN GET/HEAD requests Follow 60d9d35be282f97b497d056875e7e6ba76aaf47b. Signed-off-by: Leonard Liubich --- pkg/services/object/get.go | 36 +--- pkg/services/object/head.go | 2 +- pkg/services/object/proto.go | 51 +++++- pkg/services/object/server.go | 177 +++++++++++--------- pkg/services/object/server_internal_test.go | 20 --- 5 files changed, 148 insertions(+), 138 deletions(-) diff --git a/pkg/services/object/get.go b/pkg/services/object/get.go index a466cdd25b..3ccf6090c1 100644 --- a/pkg/services/object/get.go +++ b/pkg/services/object/get.go @@ -21,7 +21,6 @@ import ( cid "github.com/nspcc-dev/neofs-sdk-go/container/id" "github.com/nspcc-dev/neofs-sdk-go/object" oid "github.com/nspcc-dev/neofs-sdk-go/object/id" - protoencoding "github.com/nspcc-dev/neofs-sdk-go/proto/encoding" protoobject "github.com/nspcc-dev/neofs-sdk-go/proto/object" iprotobuf "github.com/nspcc-dev/neofs-sdk-go/proto/protobuf" "github.com/nspcc-dev/neofs-sdk-go/proto/protobuf/protoscan" @@ -85,7 +84,7 @@ func callRange(ctx context.Context, conn *grpc.ClientConn, request any) (grpc.Cl // - [apistatus.ErrObjectNotFound] on 404 status // - nil on other API statuses // - any other transport/protocol error otherwise -func (x *getProxyContext) continueWithConn(ctx context.Context, req *protoobject.GetRequest, conn *grpc.ClientConn) error { +func (x *getProxyContext) continueWithConn(ctx context.Context, req mem.Buffer, body *protoobject.GetRequest_Body, conn *grpc.ClientConn) error { stream, err := callGet(ctx, conn, req) if err != nil { return err @@ -96,12 +95,12 @@ func (x *getProxyContext) continueWithConn(ctx context.Context, req *protoobject var respBuf mem.BufferSlice if err = stream.RecvMsg(&respBuf); err != nil { if errors.Is(err, io.EOF) { - return x.validateEOF(req.Body, prog) + return x.validateEOF(body, prog) } return fmt.Errorf("reading the response failed: %w", err) } - fin, sent, err := x.handleGetResponse(ctx, req.Body.Raw, &prog, respBuf) + fin, sent, err := x.handleGetResponse(ctx, body.Raw, &prog, respBuf) if !sent { respBuf.Free() } @@ -925,35 +924,16 @@ func (s *Server) makeGetECPartRequest(needSign bool, remoteServerAPIVersion vers bodyLen := protoobject.CalculateGetRequestBodyLength(false, rngOff, rngLen, payloadOnly, nil, nil) - metaHdrLen := calculateRequestMetaHeaderLen(remoteServerAPIVersion, 1, xHdrs) - - reqLen := protoencoding.CalculateRequestBodyWithMetaHeaderLength(bodyLen, metaHdrLen) + var sigCount int if needSign { - reqLen += calculateRequestVerificationHeaderFieldLen(remoteServerAPIVersion) + sigCount = calculateSignatureCountForAPIVersion(remoteServerAPIVersion) } - bufItem := defaultGRPCBufferPool.Get(reqLen) - buf := *bufItem - - // body - off := protoencoding.WriteRequestBodyTagAndLength(buf, bodyLen) - off += protoobject.WriteGetRequestBody(buf[off:], cnr, parent, false, rngOff, rngLen, payloadOnly, nil, nil) - bodySlice := buf[off-bodyLen : off] - - // meta header - off += writeRequestMetaHeaderToRequest(buf[off:], remoteServerAPIVersion, 1, xHdrs) - - if !needSign { - return bufItem, nil - } - - // verification header - err := s.writeRequestSignatures(buf, off, bodySlice, buf[off-metaHdrLen:off], remoteServerAPIVersion) - if err != nil { - return nil, err + writeBodyFn := func(buf []byte) int { + return protoobject.WriteGetRequestBody(buf, cnr, parent, false, rngOff, rngLen, payloadOnly, nil, nil) } - return bufItem, nil + return s.makeLocalRequest(sigCount, remoteServerAPIVersion, bodyLen, writeBodyFn, xHdrs) } func (x *getECTransport) makeGetECPartRangeRequest(needSign bool, remoteServerAPIVersion version.Version, partInfo iec.PartInfo, off, ln uint64) (*[]byte, error) { diff --git a/pkg/services/object/head.go b/pkg/services/object/head.go index ec07551611..a3702d198c 100644 --- a/pkg/services/object/head.go +++ b/pkg/services/object/head.go @@ -28,7 +28,7 @@ func callHead(ctx context.Context, conn *grpc.ClientConn, req any) (mem.BufferSl // - (nil, nil, [apistatus.ErrObjectNotFound]) on 404 status // - (buffered response, nil, nil) on other API statuses // - (nil, nil, err) on any transport err -func getHeaderFromRemoteNode(ctx context.Context, conn *grpc.ClientConn, req *protoobject.HeadRequest, reqOID oid.ID) (mem.BufferSlice, iprotobuf.BuffersSlice, error) { +func getHeaderFromRemoteNode(ctx context.Context, conn *grpc.ClientConn, req mem.Buffer, reqOID oid.ID) (mem.BufferSlice, iprotobuf.BuffersSlice, error) { respBuf, err := callHead(ctx, conn, req) if err != nil { return nil, iprotobuf.BuffersSlice{}, err diff --git a/pkg/services/object/proto.go b/pkg/services/object/proto.go index 3368768b12..2045e80db8 100644 --- a/pkg/services/object/proto.go +++ b/pkg/services/object/proto.go @@ -288,18 +288,19 @@ func signECDSAWithSHA512(privKey ecdsa.PrivateKey, data []byte) ([]byte, error) return sig, nil } -func calculateRequestVerificationHeaderFieldLen(apiVersion version.Version) int { - var sigCount int - +func calculateSignatureCountForAPIVersion(apiVersion version.Version) int { switch apiVersion.Compare(version.New(2, 25)) { default: - sigCount = 1 + return 1 case -1: - sigCount = 3 + return 3 case 0: - sigCount = 2 + return 2 } +} +func calculateRequestVerificationHeaderFieldLen(apiVersion version.Version) int { + sigCount := calculateSignatureCountForAPIVersion(apiVersion) return protoencoding.CalculateRequestVerificationHeaderFieldLength(sigCount * verificationHeaderECDSAWithSHA512SignatureLen) } @@ -360,3 +361,41 @@ func encodeRequestProtobuf(body protoencoding.Message, metaHdr protoencoding.Mes return bufItem } + +func (s *Server) makeLocalRequestFromBody(sigCount int, remoteServerAPIVersion version.Version, body protoencoding.Message) (*[]byte, error) { + bodyLen := body.MarshaledSize() + writeBodyFn := protoencoding.WriteStablyMarshalledMessageFunc(body) + return s.makeLocalRequest(sigCount, remoteServerAPIVersion, bodyLen, writeBodyFn, nil) +} + +func (s *Server) makeLocalRequest(sigCount int, remoteServerAPIVersion version.Version, bodyLen int, writeBodyFn protoencoding.WriteMessageFunc, xHeaders []string) (*[]byte, error) { + metaHdrLen := calculateRequestMetaHeaderLen(remoteServerAPIVersion, 1, xHeaders) + + reqLen := protoencoding.CalculateRequestBodyWithMetaHeaderLength(bodyLen, metaHdrLen) + if sigCount > 0 { + reqLen += protoencoding.CalculateRequestVerificationHeaderFieldLength(sigCount * verificationHeaderECDSAWithSHA512SignatureLen) + } + + bufItem := defaultGRPCBufferPool.Get(reqLen) + buf := *bufItem + + // body + off := protoencoding.WriteRequestBodyTagAndLength(buf, bodyLen) + off += writeBodyFn(buf[off:]) + bodySlice := buf[off-bodyLen : off] + + // meta header + off += writeRequestMetaHeaderToRequest(buf[off:], remoteServerAPIVersion, 1, xHeaders) + + if sigCount == 0 { + return bufItem, nil + } + + // verification header + err := s.writeRequestSignatures(buf, off, bodySlice, buf[off-metaHdrLen:off], remoteServerAPIVersion) + if err != nil { + return nil, err + } + + return bufItem, nil +} diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index 7476cac9db..78292d6a52 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -703,7 +703,7 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) } var resp protoobject.HeadResponse - p, err := convertHeadPrm(s.signer, reqInfo.Container, req, &resp, cnrID, objID, reqMD, s.log) + p, err := convertHeadPrm(reqInfo.Container, req, &resp, cnrID, objID, reqMD) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -711,6 +711,44 @@ func (s *Server) HeadBuffered(ctx context.Context, req *protoobject.HeadRequest) return s.makeStatusHeadResponse(req, err, needSignResp) } + var remoteReqBufs [4]*[]byte // index corresponds to signature count + defer func() { + for i := range remoteReqBufs { + if remoteReqBufs[i] != nil { + defaultGRPCBufferPool.Put(remoteReqBufs[i]) + } + } + }() + + p.SetTransportFunc(func(ctx context.Context, c clientcore.MultiAddressClient) (mem.BufferSlice, iprotobuf.BuffersSlice, error) { + cv := c.APIVersion() + apiVersion := chooseAPIVersionForNewRequest(version.New(cv.GetMajor(), cv.GetMinor())) + + var sigCount int + if !clientcore.IsMutuallyAuthenticated(c) { + sigCount = calculateSignatureCountForAPIVersion(apiVersion) + } + + if remoteReqBufs[sigCount] == nil { + var err error + remoteReqBufs[sigCount], err = s.makeLocalRequestFromBody(sigCount, apiVersion, body) + if err != nil { + return nil, iprotobuf.BuffersSlice{}, fmt.Errorf("make request (signature count = %d): %w", sigCount, err) + } + } + + var respBuf mem.BufferSlice + var hdr iprotobuf.BuffersSlice + return respBuf, hdr, c.ForAnyGRPCConn(ctx, func(ctx context.Context, conn *grpc.ClientConn) error { + var err error + respBuf, hdr, err = getHeaderFromRemoteNode(ctx, conn, mem.SliceBuffer(*remoteReqBufs[sigCount]), objID) + if err != nil { + s.log.Debug("failed to get object header from remote node", zap.Error(err)) + } + return err + }) + }) + var forwardResp mem.BufferSlice var forwardReqBufItem *[]byte @@ -829,7 +867,7 @@ func (x *headResponse) WriteHeader(hdr *object.Object) error { // converts original request into parameters accepted by the internal handler. // Note that the response is untouched within this call. -func convertHeadPrm(signer ecdsa.PrivateKey, cnr container.Container, req *protoobject.HeadRequest, resp *protoobject.HeadResponse, cnrID cid.ID, objID oid.ID, reqMD requestMetadata, log *zap.Logger) (getsvc.HeadPrm, error) { +func convertHeadPrm(cnr container.Container, req *protoobject.HeadRequest, resp *protoobject.HeadResponse, cnrID cid.ID, objID oid.ID, reqMD requestMetadata) (getsvc.HeadPrm, error) { cp := objutil.CommonPrmFromRequest(reqMD.ttl, reqMD.xHeaders, reqMD.tokens) var p getsvc.HeadPrm @@ -849,38 +887,6 @@ func convertHeadPrm(signer ecdsa.PrivateKey, cnr container.Container, req *proto return getsvc.HeadPrm{}, errors.New("missing meta header") } - var updatedRequest bool - - p.SetTransportFunc(func(ctx context.Context, c clientcore.MultiAddressClient) (mem.BufferSlice, iprotobuf.BuffersSlice, error) { - if !updatedRequest { - req = &protoobject.HeadRequest{ - Body: req.Body, - MetaHeader: &protosession.RequestMetaHeader{ - Version: c.APIVersion(), - Ttl: 1, - }, - } - if shouldSignOutgoingRequest(c, req) { - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) - if err != nil { - return nil, iprotobuf.BuffersSlice{}, err - } - } - updatedRequest = true - } - - var respBuf mem.BufferSlice - var hdr iprotobuf.BuffersSlice - return respBuf, hdr, c.ForAnyGRPCConn(ctx, func(ctx context.Context, conn *grpc.ClientConn) error { - var err error - respBuf, hdr, err = getHeaderFromRemoteNode(ctx, conn, req, objID) - if err != nil { - log.Debug("failed to get object header from remote node", zap.Error(err)) - } - return err - }) - }) return p, nil } @@ -1075,7 +1081,7 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ recheckEACL = true } - p, err := convertGetPrm(s.signer, reqInfo.Container, req, &getStream{ + respStream := &getStream{ base: gStream, srv: s, reqCID: cnrID, @@ -1086,7 +1092,9 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ payloadOnly: req.GetBody().GetPayloadOnly(), returnVersionInResponse: util.NeedVersionInResponse(req.MetaHeader), sendECPartIndInResponse: sendECPartIdxInResponse(req), - }, cnrID, objID, reqMD, s.log) + } + + p, err := convertGetPrm(reqInfo.Container, req, respStream, cnrID, objID, reqMD) if err != nil { if !errors.Is(err, apistatus.Error) { err = newBadRequestError(err.Error()) // defer @@ -1094,6 +1102,55 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ return s.sendStatusGetResponse(req, gStream, err, needSignResp) } + proxyCtx := getProxyContext{ + respStream: respStream, + suppressInit: body.GetPayloadOnly(), + } + if body.GetExtendedRange() != nil { + proxyCtx.resolveRange = p.ResolveRange + } + + var remoteReqBufs [4]*[]byte // index corresponds to signature count + defer func() { + for i := range remoteReqBufs { + if remoteReqBufs[i] != nil { + defaultGRPCBufferPool.Put(remoteReqBufs[i]) + } + } + }() + + p.SetTransportFunc(func(ctx context.Context, c clientcore.MultiAddressClient) error { + cv := c.APIVersion() + apiVersion := chooseAPIVersionForNewRequest(version.New(cv.GetMajor(), cv.GetMinor())) + + var sigCount int + if !clientcore.IsMutuallyAuthenticated(c) { + sigCount = calculateSignatureCountForAPIVersion(apiVersion) + } + + if remoteReqBufs[sigCount] == nil { + body := body + if proxyCtx.suppressInit { + body = proto.Clone(body).(*protoobject.GetRequest_Body) + body.PayloadOnly = false + } + + var err error + remoteReqBufs[sigCount], err = s.makeLocalRequestFromBody(sigCount, apiVersion, body) + if err != nil { + return fmt.Errorf("make request (signature count = %d): %w", sigCount, err) + } + } + + return c.ForAnyGRPCConn(ctx, func(ctx context.Context, conn *grpc.ClientConn) error { + err := proxyCtx.continueWithConn(ctx, mem.SliceBuffer(*remoteReqBufs[sigCount]), body, conn) + if err != nil { + s.log.Debug("failed to get object from remote node", zap.Error(err)) + } + return err + }) + }) + var forwardReqBufItem *[]byte defer func() { if forwardReqBufItem != nil { @@ -1368,7 +1425,7 @@ func (s *Server) copyGetStream(gStream grpc.ServerStream, hdrRespBuf *iprotobuf. // converts original request into parameters accepted by the internal handler. // Note that the stream is untouched within this call, errors are not reported // into it. -func convertGetPrm(signer ecdsa.PrivateKey, cnr container.Container, req *protoobject.GetRequest, stream *getStream, cnrID cid.ID, objID oid.ID, reqMD requestMetadata, log *zap.Logger) (getsvc.Prm, error) { +func convertGetPrm(cnr container.Container, req *protoobject.GetRequest, stream *getStream, cnrID cid.ID, objID oid.ID, reqMD requestMetadata) (getsvc.Prm, error) { body := req.GetBody() rng := body.GetRange() @@ -1435,47 +1492,6 @@ func convertGetPrm(signer ecdsa.PrivateKey, cnr container.Container, req *protoo return getsvc.Prm{}, errors.New("missing meta header") } - proxyCtx := getProxyContext{ - respStream: stream, - suppressInit: body.GetPayloadOnly(), - } - if extendedRange != nil { - proxyCtx.resolveRange = p.ResolveRange - } - - var updatedRequest bool - - p.SetTransportFunc(func(ctx context.Context, c clientcore.MultiAddressClient) error { - if !updatedRequest { - req = &protoobject.GetRequest{ - Body: req.Body, - MetaHeader: &protosession.RequestMetaHeader{ - Version: c.APIVersion(), - Ttl: 1, - }, - } - if proxyCtx.suppressInit { - req.Body = proto.Clone(req.Body).(*protoobject.GetRequest_Body) - req.Body.PayloadOnly = false - } - if shouldSignOutgoingRequest(c, req) { - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) - if err != nil { - return err - } - } - updatedRequest = true - } - - return c.ForAnyGRPCConn(ctx, func(ctx context.Context, conn *grpc.ClientConn) error { - err := proxyCtx.continueWithConn(ctx, req, conn) - if err != nil { - log.Debug("failed to get object from remote node", zap.Error(err)) - } - return err - }) - }) return p, nil } @@ -2510,11 +2526,6 @@ func needSignGetResponse(req util.Request) bool { return util.VersionLE(req, 2, 17) } -func shouldSignOutgoingRequest(c any, req util.Request) bool { - meta := req.GetMetaHeader() - return meta == nil || meta.GetTtl() != 1 || !clientcore.IsMutuallyAuthenticated(c) -} - func checkHeaderProtobufAgainstID(buffers iprotobuf.BuffersSlice, id oid.ID, ordered bool) error { b := buffers.ReadOnlyData() if !ordered { diff --git a/pkg/services/object/server_internal_test.go b/pkg/services/object/server_internal_test.go index 4c440824df..c75de452a3 100644 --- a/pkg/services/object/server_internal_test.go +++ b/pkg/services/object/server_internal_test.go @@ -5,29 +5,9 @@ import ( iec "github.com/nspcc-dev/neofs-node/internal/ec" "github.com/nspcc-dev/neofs-sdk-go/netmap" - protoobject "github.com/nspcc-dev/neofs-sdk-go/proto/object" - protosession "github.com/nspcc-dev/neofs-sdk-go/proto/session" "github.com/stretchr/testify/require" ) -type mutuallyAuthenticatedClient bool - -func (x mutuallyAuthenticatedClient) IsMutuallyAuthenticated() bool { - return bool(x) -} - -func TestShouldSignOutgoingRequest(t *testing.T) { - requestWithTTL := func(ttl uint32) *protoobject.GetRequest { - return &protoobject.GetRequest{MetaHeader: &protosession.RequestMetaHeader{Ttl: ttl}} - } - - require.False(t, shouldSignOutgoingRequest(mutuallyAuthenticatedClient(true), requestWithTTL(1))) - require.True(t, shouldSignOutgoingRequest(mutuallyAuthenticatedClient(true), requestWithTTL(2))) - require.True(t, shouldSignOutgoingRequest(mutuallyAuthenticatedClient(true), new(protoobject.GetRequest))) - require.True(t, shouldSignOutgoingRequest(mutuallyAuthenticatedClient(false), requestWithTTL(1))) - require.True(t, shouldSignOutgoingRequest(struct{}{}, requestWithTTL(1))) -} - func TestIterateSearchableContainerNodes(t *testing.T) { for _, tc := range []struct { name string