From 158f04c12c977dd6517c6bdd67cb299142e2d9cd Mon Sep 17 00:00:00 2001 From: Roman Khimov Date: Fri, 4 Sep 2026 17:26:47 +0300 Subject: [PATCH] object: drop Range API support, fix #4136 Leave stub at the node level, drop everything related otherwise, it's deprecated since API 2.23 and no longer required for applications. Signed-off-by: Roman Khimov --- cmd/neofs-cli/modules/acl/extended/create.go | 2 +- cmd/neofs-cli/modules/object/get.go | 94 +++++++ cmd/neofs-cli/modules/object/range.go | 229 ------------------ cmd/neofs-cli/modules/object/root.go | 2 - cmd/neofs-cli/modules/object/util.go | 4 - .../modules/object/util_session_v2.go | 2 - cmd/neofs-cli/modules/session/create_v2.go | 4 +- cmd/neofs-cli/modules/util/acl.go | 5 +- cmd/neofs-node/object.go | 4 - .../neofs-cli_acl_extended_create.md | 2 +- docs/cli-commands/neofs-cli_object.md | 1 - docs/cli-commands/neofs-cli_object_range.md | 43 ---- pkg/core/client/client.go | 1 - pkg/metrics/object.go | 17 -- pkg/network/cache/clients.go | 10 - pkg/services/object/acl/eacl/v2/headers.go | 1 - pkg/services/object/acl/v2/service.go | 6 - pkg/services/object/acl/v2/service_test.go | 20 -- pkg/services/object/acl/v2/util.go | 5 +- pkg/services/object/acl/v2/util_test.go | 3 - pkg/services/object/get.go | 13 - pkg/services/object/get/ec.go | 11 - pkg/services/object/get/exec.go | 34 +-- pkg/services/object/get/get.go | 70 +----- pkg/services/object/get/get_test.go | 221 +---------------- pkg/services/object/get/prm.go | 43 +--- pkg/services/object/get/util.go | 166 +------------ pkg/services/object/proto.go | 4 - pkg/services/object/range.go | 179 -------------- pkg/services/object/server.go | 202 +-------------- pkg/services/object/server_test.go | 10 - 31 files changed, 119 insertions(+), 1289 deletions(-) delete mode 100644 cmd/neofs-cli/modules/object/range.go delete mode 100644 docs/cli-commands/neofs-cli_object_range.md delete mode 100644 pkg/services/object/range.go diff --git a/cmd/neofs-cli/modules/acl/extended/create.go b/cmd/neofs-cli/modules/acl/extended/create.go index 07c3daf552..c2de895ccd 100644 --- a/cmd/neofs-cli/modules/acl/extended/create.go +++ b/cmd/neofs-cli/modules/acl/extended/create.go @@ -24,7 +24,7 @@ Rule consist of these blocks: [ ...] [ .. Action is 'allow' or 'deny'. -Operation is an object service verb: 'get', 'head', 'put', 'search', 'delete', or 'getrange'. +Operation is an object service verb: 'get', 'head', 'put', 'search' or 'delete'. Filter consists of : Typ is 'obj' for object applied filter or 'req' for request applied filter. diff --git a/cmd/neofs-cli/modules/object/get.go b/cmd/neofs-cli/modules/object/get.go index 876cf220fa..f6d417bae2 100644 --- a/cmd/neofs-cli/modules/object/get.go +++ b/cmd/neofs-cli/modules/object/get.go @@ -1,6 +1,8 @@ package object import ( + "bytes" + "errors" "fmt" "io" "os" @@ -21,6 +23,10 @@ import ( ) const ( + rangeFlag = "range" + rangeFlagUsage = "Range to take data from in the form offset:length" + rangeSep = ":" + payloadOnlyFlag = "payload-only" extendedRangeFlag = "extended-range" ) @@ -269,3 +275,91 @@ func payloadReadSize(payloadSize uint64, ranges []*object.Range, first, last *ui return int64(payloadSize) } } + +func printSplitInfoErr(cmd *cobra.Command, err error) (bool, error) { + errSplitInfo, ok := errors.AsType[*object.SplitInfoError](err) + + if ok { + cmd.PrintErrln("Object is complex, split information received.") + if err := printSplitInfo(cmd, errSplitInfo.SplitInfo()); err != nil { + return false, err + } + } + + return ok, nil +} + +func printSplitInfo(cmd *cobra.Command, info *object.SplitInfo) error { + bs, err := marshalSplitInfo(cmd, info) + if err != nil { + return fmt.Errorf("can't marshal split info: %w", err) + } + + cmd.Println(string(bs)) + return nil +} + +func marshalSplitInfo(cmd *cobra.Command, info *object.SplitInfo) ([]byte, error) { + toJSON, _ := cmd.Flags().GetBool(commonflags.JSON) + toProto, _ := cmd.Flags().GetBool("proto") + switch { + case toJSON && toProto: + return nil, errors.New("'--json' and '--proto' flags are mutually exclusive") + case toJSON: + return info.MarshalJSON() + case toProto: + return info.Marshal(), nil + default: + b := bytes.NewBuffer(nil) + if splitID := info.SplitID(); splitID != nil { + b.WriteString("Split ID: " + splitID.String() + "\n") + } + if link := info.GetLink(); !link.IsZero() { + b.WriteString("Linking object: " + link.String() + "\n") + } + if first := info.GetFirstPart(); !first.IsZero() { + b.WriteString("First object: " + first.String() + "\n") + } + if last := info.GetLastPart(); !last.IsZero() { + b.WriteString("Last object: " + last.String() + "\n") + } + return b.Bytes(), nil + } +} + +func getRangeList(cmd *cobra.Command) ([]*object.Range, error) { + v := cmd.Flag("range").Value.String() + if len(v) == 0 { + return nil, nil + } + vs := strings.Split(v, ",") + rs := make([]*object.Range, len(vs)) + for i := range vs { + offString, lenString, found := strings.Cut(vs[i], rangeSep) + if !found { + return nil, fmt.Errorf("invalid range specifier: %s", vs[i]) + } + + offset, err := strconv.ParseUint(offString, 10, 64) + if err != nil { + return nil, fmt.Errorf("invalid '%s' range offset specifier: %w", vs[i], err) + } + length, err := strconv.ParseUint(lenString, 10, 64) + if err != nil { + return nil, fmt.Errorf("invalid '%s' range length specifier: %w", vs[i], err) + } + + if length == 0 { + if offset != 0 { + return nil, fmt.Errorf("invalid '%s' range: zero length with non-zero offset", vs[i]) + } + } else if offset+length <= offset { + return nil, fmt.Errorf("invalid '%s' range: uint64 overflow", vs[i]) + } + + rs[i] = object.NewRange() + rs[i].SetOffset(offset) + rs[i].SetLength(length) + } + return rs, nil +} diff --git a/cmd/neofs-cli/modules/object/range.go b/cmd/neofs-cli/modules/object/range.go deleted file mode 100644 index 085dea7438..0000000000 --- a/cmd/neofs-cli/modules/object/range.go +++ /dev/null @@ -1,229 +0,0 @@ -package object - -import ( - "bytes" - "errors" - "fmt" - "io" - "os" - "strconv" - "strings" - - internalclient "github.com/nspcc-dev/neofs-node/cmd/neofs-cli/internal/client" - "github.com/nspcc-dev/neofs-node/cmd/neofs-cli/internal/commonflags" - "github.com/nspcc-dev/neofs-node/cmd/neofs-cli/internal/key" - "github.com/nspcc-dev/neofs-sdk-go/client" - 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" - "github.com/nspcc-dev/neofs-sdk-go/user" - "github.com/spf13/cobra" -) - -const ( - rangeSep = ":" - rangeFlag = "range" - rangeFlagUsage = "Range to take data from in the form offset:length" -) - -var objectRangeCmd = &cobra.Command{ - Use: "range", - Short: "Get payload range data of an object", - Long: "Get payload range data of an object", - Args: cobra.NoArgs, - RunE: getObjectRange, -} - -func initObjectRangeCmd() { - commonflags.Init(objectRangeCmd) - initFlagSession(objectRangeCmd, "RANGE") - - flags := objectRangeCmd.Flags() - - flags.String(commonflags.CIDFlag, "", commonflags.CIDFlagUsage) - _ = objectRangeCmd.MarkFlagRequired(commonflags.CIDFlag) - - flags.String(commonflags.OIDFlag, "", commonflags.OIDFlagUsage) - _ = objectRangeCmd.MarkFlagRequired(commonflags.OIDFlag) - - flags.String("range", "", "Range to take data from in the form offset:length") - flags.String(fileFlag, "", "File to write object payload to. Default: stdout.") - flags.Bool(rawFlag, false, rawFlagDesc) -} - -func getObjectRange(cmd *cobra.Command, _ []string) error { - var cnr cid.ID - var obj oid.ID - - _, err := readObjectAddress(cmd, &cnr, &obj) - if err != nil { - return err - } - - ranges, err := getRangeList(cmd) - if err != nil { - return err - } - - if len(ranges) != 1 { - return fmt.Errorf("exactly one range must be specified, got: %d", len(ranges)) - } - - var out io.Writer - - filename := cmd.Flag(fileFlag).Value.String() - if filename == "" { - out = os.Stdout - } else { - f, err := openFileForPayload(filename) - if err != nil { - return err - } - - defer f.Close() - - out = f - } - - pk, err := key.GetOrGenerate(cmd) - if err != nil { - return err - } - - ctx, cancel := commonflags.GetCommandContext(cmd) - defer cancel() - - cli, err := internalclient.GetSDKClientByFlag(ctx, commonflags.RPC) - if err != nil { - return err - } - defer cli.Close() - - var prm client.PrmObjectRange - err = Prepare(cmd, &prm) - if err != nil { - return err - } - err = readSession(cmd, &prm, pk, cnr, obj) - if err != nil { - return err - } - - raw, _ := cmd.Flags().GetBool(rawFlag) - if raw { - prm.MarkRaw() - } - - //nolint:staticcheck - rdr, err := cli.ObjectRangeInit(ctx, cnr, obj, ranges[0].GetOffset(), ranges[0].GetLength(), user.NewAutoIDSigner(*pk), prm) - if err != nil { - err = fmt.Errorf("init payload reading: %w", err) - } else { - if _, err = io.Copy(out, rdr); err != nil { - err = fmt.Errorf("copy payload: %w", err) - } - } - if err != nil { - if ok, err := printSplitInfoErr(cmd, err); ok { - return err - } - - if err != nil { - return fmt.Errorf("can't get object payload range: %w", err) - } - } - - if filename != "" { - cmd.Printf("[%s] Payload successfully saved\n", filename) - } - - return nil -} - -func printSplitInfoErr(cmd *cobra.Command, err error) (bool, error) { - errSplitInfo, ok := errors.AsType[*object.SplitInfoError](err) - - if ok { - cmd.PrintErrln("Object is complex, split information received.") - if err := printSplitInfo(cmd, errSplitInfo.SplitInfo()); err != nil { - return false, err - } - } - - return ok, nil -} - -func printSplitInfo(cmd *cobra.Command, info *object.SplitInfo) error { - bs, err := marshalSplitInfo(cmd, info) - if err != nil { - return fmt.Errorf("can't marshal split info: %w", err) - } - - cmd.Println(string(bs)) - return nil -} - -func marshalSplitInfo(cmd *cobra.Command, info *object.SplitInfo) ([]byte, error) { - toJSON, _ := cmd.Flags().GetBool(commonflags.JSON) - toProto, _ := cmd.Flags().GetBool("proto") - switch { - case toJSON && toProto: - return nil, errors.New("'--json' and '--proto' flags are mutually exclusive") - case toJSON: - return info.MarshalJSON() - case toProto: - return info.Marshal(), nil - default: - b := bytes.NewBuffer(nil) - if splitID := info.SplitID(); splitID != nil { - b.WriteString("Split ID: " + splitID.String() + "\n") - } - if link := info.GetLink(); !link.IsZero() { - b.WriteString("Linking object: " + link.String() + "\n") - } - if first := info.GetFirstPart(); !first.IsZero() { - b.WriteString("First object: " + first.String() + "\n") - } - if last := info.GetLastPart(); !last.IsZero() { - b.WriteString("Last object: " + last.String() + "\n") - } - return b.Bytes(), nil - } -} - -func getRangeList(cmd *cobra.Command) ([]*object.Range, error) { - v := cmd.Flag("range").Value.String() - if len(v) == 0 { - return nil, nil - } - vs := strings.Split(v, ",") - rs := make([]*object.Range, len(vs)) - for i := range vs { - offString, lenString, found := strings.Cut(vs[i], rangeSep) - if !found { - return nil, fmt.Errorf("invalid range specifier: %s", vs[i]) - } - - offset, err := strconv.ParseUint(offString, 10, 64) - if err != nil { - return nil, fmt.Errorf("invalid '%s' range offset specifier: %w", vs[i], err) - } - length, err := strconv.ParseUint(lenString, 10, 64) - if err != nil { - return nil, fmt.Errorf("invalid '%s' range length specifier: %w", vs[i], err) - } - - if length == 0 { - if offset != 0 { - return nil, fmt.Errorf("invalid '%s' range: zero length with non-zero offset", vs[i]) - } - } else if offset+length <= offset { - return nil, fmt.Errorf("invalid '%s' range: uint64 overflow", vs[i]) - } - - rs[i] = object.NewRange() - rs[i].SetOffset(offset) - rs[i].SetLength(length) - } - return rs, nil -} diff --git a/cmd/neofs-cli/modules/object/root.go b/cmd/neofs-cli/modules/object/root.go index d49f333241..530040e16d 100644 --- a/cmd/neofs-cli/modules/object/root.go +++ b/cmd/neofs-cli/modules/object/root.go @@ -26,7 +26,6 @@ func init() { objectSearchCmd, searchV2Cmd, objectHeadCmd, - objectRangeCmd, objectLockCmd} Cmd.AddCommand(objectNodesCmd) @@ -42,7 +41,6 @@ func init() { initObjectGetCmd() initObjectSearchCmd() initObjectHeadCmd() - initObjectRangeCmd() initCommandObjectLock() initObjectNodesCmd() } diff --git a/cmd/neofs-cli/modules/object/util.go b/cmd/neofs-cli/modules/object/util.go index c72151c2bd..ffc25ef89c 100644 --- a/cmd/neofs-cli/modules/object/util.go +++ b/cmd/neofs-cli/modules/object/util.go @@ -197,8 +197,6 @@ func getSession(cmd *cobra.Command) (*session.Object, error) { // *internal.GetObjectPrm // *internal.HeadObjectPrm // *internal.SearchObjectsPrm -// *internal.PayloadRangePrm -// *internal.HashPayloadRangesPrm func _readVerifiedSession(cmd *cobra.Command, dst SessionPrm, key *ecdsa.PrivateKey, cnr cid.ID, obj *oid.ID) error { if tokV2 := tryReadSessionV2(cmd); tokV2 != nil { err := attachVerifiedSessionV2(cmd, tokV2, dst, key, cnr) @@ -222,8 +220,6 @@ func _readVerifiedSession(cmd *cobra.Command, dst SessionPrm, key *ecdsa.Private cmdVerb = session.VerbObjectHead case *client.PrmObjectSearch: cmdVerb = session.VerbObjectSearch - case *client.PrmObjectRange: - cmdVerb = session.VerbObjectRange } tok, err := getVerifiedSession(cmd, cmdVerb, key, cnr) diff --git a/cmd/neofs-cli/modules/object/util_session_v2.go b/cmd/neofs-cli/modules/object/util_session_v2.go index f570b1256e..d18e9dc023 100644 --- a/cmd/neofs-cli/modules/object/util_session_v2.go +++ b/cmd/neofs-cli/modules/object/util_session_v2.go @@ -77,8 +77,6 @@ func attachVerifiedSessionV2(cmd *cobra.Command, tok *session.Token, dst Session cmdVerb = session.VerbObjectHead case *client.PrmObjectSearch: cmdVerb = session.VerbObjectSearch - case *client.PrmObjectRange: - cmdVerb = session.VerbObjectRange } err := verifySessionV2(cmd, tok, cmdVerb, key, cnr) diff --git a/cmd/neofs-cli/modules/session/create_v2.go b/cmd/neofs-cli/modules/session/create_v2.go index 89fac30085..3494e9ad6f 100644 --- a/cmd/neofs-cli/modules/session/create_v2.go +++ b/cmd/neofs-cli/modules/session/create_v2.go @@ -366,8 +366,6 @@ func parseVerbs(verbsStr string) ([]session.Verb, error) { verb = session.VerbObjectSearch case "DELETE", "OBJECTDELETE": verb = session.VerbObjectDelete - case "RANGE", "OBJECTRANGE": - verb = session.VerbObjectRange case "CONTAINERSET", "CONTAINERSETACL", "CONTAINER_SET", "CONTAINER_SET_ACL": verb = session.VerbContainerSetEACL case "CONTAINERPUT", "CONTAINER_PUT": @@ -375,7 +373,7 @@ func parseVerbs(verbsStr string) ([]session.Verb, error) { case "CONTAINERDELETE", "CONTAINER_DELETE": verb = session.VerbContainerDelete default: - return nil, fmt.Errorf("unknown verb: %s (supported: GET,PUT,HEAD,SEARCH,DELETE,RANGE,CONTAINERSET,CONTAINERPUT,CONTAINERDELETE)", verbStr) + return nil, fmt.Errorf("unknown verb: %s (supported: GET,PUT,HEAD,SEARCH,DELETE,CONTAINERSET,CONTAINERPUT,CONTAINERDELETE)", verbStr) } verbs = append(verbs, verb) diff --git a/cmd/neofs-cli/modules/util/acl.go b/cmd/neofs-cli/modules/util/acl.go index aacc61b79c..c70b17a7ea 100644 --- a/cmd/neofs-cli/modules/util/acl.go +++ b/cmd/neofs-cli/modules/util/acl.go @@ -20,11 +20,10 @@ import ( func PrettyPrintTableBACL(cmd *cobra.Command, bacl *acl.Basic) { // Header w := tabwriter.NewWriter(cmd.OutOrStdout(), 1, 4, 4, ' ', 0) - fmt.Fprintln(w, "\tRange\tSearch\tDelete\tPut\tHead\tGet") + fmt.Fprintln(w, "\tSearch\tDelete\tPut\tHead\tGet") // Bits bits := []string{ boolToString(bacl.Sticky()) + " " + boolToString(!bacl.Extendable()), - getRoleBitsForOperation(bacl, acl.OpObjectRange), getRoleBitsForOperation(bacl, acl.OpObjectSearch), getRoleBitsForOperation(bacl, acl.OpObjectDelete), getRoleBitsForOperation(bacl, acl.OpObjectPut), getRoleBitsForOperation(bacl, acl.OpObjectHead), getRoleBitsForOperation(bacl, acl.OpObjectGet), @@ -32,7 +31,7 @@ func PrettyPrintTableBACL(cmd *cobra.Command, bacl *acl.Basic) { fmt.Fprintln(w, strings.Join(bits, "\t")) // Footer footer := []string{"X F"} - for range 6 { + for range 5 { footer = append(footer, "U S O B") } fmt.Fprintln(w, strings.Join(footer, "\t")) diff --git a/cmd/neofs-node/object.go b/cmd/neofs-node/object.go index 8f585497dd..4a86003783 100644 --- a/cmd/neofs-node/object.go +++ b/cmd/neofs-node/object.go @@ -125,10 +125,6 @@ func (s *objectSvc) Delete(ctx context.Context, prm deletesvc.Prm) error { return s.delete.Delete(ctx, prm) } -func (s *objectSvc) GetRange(ctx context.Context, prm getsvc.RangePrm) error { - return s.get.GetRange(ctx, prm) -} - type delNetInfo struct { netmapcore.State tsLifetime uint64 diff --git a/docs/cli-commands/neofs-cli_acl_extended_create.md b/docs/cli-commands/neofs-cli_acl_extended_create.md index d00520eb6c..0df679858a 100644 --- a/docs/cli-commands/neofs-cli_acl_extended_create.md +++ b/docs/cli-commands/neofs-cli_acl_extended_create.md @@ -10,7 +10,7 @@ Rule consist of these blocks: [ ...] [ .. Action is 'allow' or 'deny'. -Operation is an object service verb: 'get', 'head', 'put', 'search', 'delete', or 'getrange'. +Operation is an object service verb: 'get', 'head', 'put', 'search' or 'delete'. Filter consists of : Typ is 'obj' for object applied filter or 'req' for request applied filter. diff --git a/docs/cli-commands/neofs-cli_object.md b/docs/cli-commands/neofs-cli_object.md index a12b28650a..47d1772063 100644 --- a/docs/cli-commands/neofs-cli_object.md +++ b/docs/cli-commands/neofs-cli_object.md @@ -28,7 +28,6 @@ Operations with Objects * [neofs-cli object lock](neofs-cli_object_lock.md) - Lock object in container * [neofs-cli object nodes](neofs-cli_object_nodes.md) - Show nodes for an object * [neofs-cli object put](neofs-cli_object_put.md) - Put object to NeoFS -* [neofs-cli object range](neofs-cli_object_range.md) - Get payload range data of an object * [neofs-cli object search](neofs-cli_object_search.md) - Search object * [neofs-cli object searchv2](neofs-cli_object_searchv2.md) - Search object (deprecated) diff --git a/docs/cli-commands/neofs-cli_object_range.md b/docs/cli-commands/neofs-cli_object_range.md deleted file mode 100644 index ff3318b9c8..0000000000 --- a/docs/cli-commands/neofs-cli_object_range.md +++ /dev/null @@ -1,43 +0,0 @@ -## neofs-cli object range - -Get payload range data of an object - -### Synopsis - -Get payload range data of an object - -``` -neofs-cli object range [flags] -``` - -### Options - -``` - --address string Address of wallet account - --bearer string File with signed JSON or binary encoded bearer token - --cid string Container ID. - --file string File to write object payload to. Default: stdout. - -g, --generate-key Generate new private key - -h, --help help for range - --oid string Object ID. - --range string Range to take data from in the form offset:length - --raw Set raw request option - -r, --rpc-endpoint string Remote node address (as 'multiaddr' or ':') - --session string Filepath to a JSON- or binary-encoded token of the object RANGE session - -t, --timeout duration Timeout for the operation (default 15s) - --ttl uint32 TTL value in request meta header (default 2) - -w, --wallet string Path to the wallet - -x, --xhdr strings Request X-Headers in form of Key=Value -``` - -### Options inherited from parent commands - -``` - -c, --config string Config file (default is $HOME/.config/neofs-cli/config.yaml) - -v, --verbose Verbose output -``` - -### SEE ALSO - -* [neofs-cli object](neofs-cli_object.md) - Operations with Objects - diff --git a/pkg/core/client/client.go b/pkg/core/client/client.go index 345edf75c9..362b7959df 100644 --- a/pkg/core/client/client.go +++ b/pkg/core/client/client.go @@ -26,7 +26,6 @@ type Client interface { ObjectHead(ctx context.Context, containerID cid.ID, objectID oid.ID, signer user.Signer, prm client.PrmObjectHead) (*object.Object, error) ObjectSearchInit(ctx context.Context, containerID cid.ID, signer user.Signer, prm client.PrmObjectSearch) (*client.ObjectListReader, error) SearchObjects(context.Context, cid.ID, object.SearchFilters, []string, string, neofscrypto.Signer, client.SearchObjectsOptions) ([]client.SearchResultItem, string, error) - ObjectRangeInit(ctx context.Context, containerID cid.ID, objectID oid.ID, offset, length uint64, signer user.Signer, prm client.PrmObjectRange) (*client.ObjectRangeReader, error) AnnounceLocalTrust(ctx context.Context, epoch uint64, trusts []reputation.Trust, prm client.PrmAnnounceLocalTrust) error AnnounceIntermediateTrust(ctx context.Context, epoch uint64, trust reputation.PeerToPeerTrust, prm client.PrmAnnounceIntermediateTrust) error } diff --git a/pkg/metrics/object.go b/pkg/metrics/object.go index ea834c1d89..886f1c09b8 100644 --- a/pkg/metrics/object.go +++ b/pkg/metrics/object.go @@ -22,14 +22,12 @@ type ( headCounter methodCount searchCounter methodCount deleteCounter methodCount - rangeCounter methodCount getDuration prometheus.Histogram putDuration prometheus.Histogram headDuration prometheus.Histogram searchDuration prometheus.Histogram deleteDuration prometheus.Histogram - rangeDuration prometheus.Histogram putPayload prometheus.Counter getPayload prometheus.Counter @@ -81,7 +79,6 @@ func newObjectServiceMetrics() objectServiceMetrics { headCounter = newMethodCallCounter("head") searchCounter = newMethodCallCounter("search") deleteCounter = newMethodCallCounter("delete") - rangeCounter = newMethodCallCounter("range") ) var ( // Request duration metrics. @@ -119,13 +116,6 @@ func newObjectServiceMetrics() objectServiceMetrics { Name: "rpc_delete_time", Help: "RPC 'delete' request handling time", }) - - rangeDuration = prometheus.NewHistogram(prometheus.HistogramOpts{ - Namespace: storageNodeNameSpace, - Subsystem: objectSubsystem, - Name: "rpc_range_time", - Help: "RPC 'range request' handling time", - }) ) var ( // Object payload metrics. @@ -168,13 +158,11 @@ func newObjectServiceMetrics() objectServiceMetrics { headCounter: headCounter, searchCounter: searchCounter, deleteCounter: deleteCounter, - rangeCounter: rangeCounter, getDuration: getDuration, putDuration: putDuration, headDuration: headDuration, searchDuration: searchDuration, deleteDuration: deleteDuration, - rangeDuration: rangeDuration, putPayload: putPayload, getPayload: getPayload, shardMetrics: shardsMetrics, @@ -188,14 +176,12 @@ func (m objectServiceMetrics) register() { m.headCounter.mustRegister() m.searchCounter.mustRegister() m.deleteCounter.mustRegister() - m.rangeCounter.mustRegister() prometheus.MustRegister(m.getDuration) prometheus.MustRegister(m.putDuration) prometheus.MustRegister(m.headDuration) prometheus.MustRegister(m.searchDuration) prometheus.MustRegister(m.deleteDuration) - prometheus.MustRegister(m.rangeDuration) prometheus.MustRegister(m.putPayload) prometheus.MustRegister(m.getPayload) @@ -223,9 +209,6 @@ func (m objectServiceMetrics) HandleOpExecResult(op stat.Method, success bool, d case stat.MethodObjectSearch, stat.MethodObjectSearchV2: // FIXME: sep counters? m.searchCounter.inc(success) m.searchDuration.Observe(d.Seconds()) - case stat.MethodObjectRange: - m.rangeCounter.inc(success) - m.rangeDuration.Observe(d.Seconds()) } } diff --git a/pkg/network/cache/clients.go b/pkg/network/cache/clients.go index c071e6eccd..e9cb874fb2 100644 --- a/pkg/network/cache/clients.go +++ b/pkg/network/cache/clients.go @@ -486,16 +486,6 @@ func (x *connections) ObjectSearchInit(ctx context.Context, cnr cid.ID, signer u }) } -func (x *connections) ObjectRangeInit(ctx context.Context, cnr cid.ID, id oid.ID, off, ln uint64, signer user.Signer, opts client.PrmObjectRange) (*client.ObjectRangeReader, error) { - var res *client.ObjectRangeReader - return res, x.forAny(ctx, func(ctx context.Context, c *client.Client) error { - var err error - //nolint:staticcheck - res, err = c.ObjectRangeInit(ctx, cnr, id, off, ln, signer, opts) - return err - }) -} - func (x *connections) AnnounceLocalTrust(ctx context.Context, epoch uint64, ts []reputation.Trust, opts client.PrmAnnounceLocalTrust) error { return x.forAny(ctx, func(ctx context.Context, c *client.Client) error { return c.AnnounceLocalTrust(ctx, epoch, ts, opts) diff --git a/pkg/services/object/acl/eacl/v2/headers.go b/pkg/services/object/acl/eacl/v2/headers.go index 0592b8106c..52af580451 100644 --- a/pkg/services/object/acl/eacl/v2/headers.go +++ b/pkg/services/object/acl/eacl/v2/headers.go @@ -136,7 +136,6 @@ func (h *cfg) readObjectHeaders(dst *headerSource) error { dst.objectHeaders = objHeaders dst.incompleteObjectHeaders = !completed case - *protoobject.GetRangeRequest, *protoobject.DeleteRequest: dst.objectHeaders = addressHeaders(h.cnr, h.obj) case *protoobject.PutRequest: diff --git a/pkg/services/object/acl/v2/service.go b/pkg/services/object/acl/v2/service.go index 7181416577..0f8cd81964 100644 --- a/pkg/services/object/acl/v2/service.go +++ b/pkg/services/object/acl/v2/service.go @@ -449,12 +449,6 @@ func (b Service) DeleteRequestToInfo(ctx context.Context, request *protoobject.D return b.findRequestInfo(ctx, request, cnr, acl.OpObjectDelete, tokens) } -// RangeRequestToInfo resolves RequestInfo from the request to check it using -// [ACLChecker]. -func (b Service) RangeRequestToInfo(ctx context.Context, request *protoobject.GetRangeRequest, cnr cid.ID, tokens common.RequestTokens) (RequestInfo, error) { - return b.findRequestInfo(ctx, request, cnr, acl.OpObjectRange, tokens) -} - var ErrSkipRequest = errors.New("skip request") // PutRequestToInfo resolves RequestInfo from the request to check it using diff --git a/pkg/services/object/acl/v2/service_test.go b/pkg/services/object/acl/v2/service_test.go index 1b70619505..41765a8ac8 100644 --- a/pkg/services/object/acl/v2/service_test.go +++ b/pkg/services/object/acl/v2/service_test.go @@ -242,26 +242,6 @@ func TestService_DeleteRequestToInfo_BearerTokenIssuer(t *testing.T) { }) } -func TestService_RangeRequestToInfo_BearerTokenIssuer(t *testing.T) { - testBearerTokenIssuer(t, (*aclsvc.Service).RangeRequestToInfo, func(t *testing.T, signer neofscrypto.Signer, cnrID cid.ID, meta *protosession.RequestMetaHeader) *protoobject.GetRangeRequest { - req := &protoobject.GetRangeRequest{ - Body: &protoobject.GetRangeRequest_Body{ - Address: &refs.Address{ - ContainerId: cnrID.ProtoMessage(), - ObjectId: oidtest.ID().ProtoMessage(), - }, - }, - MetaHeader: meta, - } - - var err error - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(signer, req, nil) - require.NoError(t, err) - - return req - }) -} - func TestService_PutRequestToInfo_BearerTokenIssuer(t *testing.T) { testBearerTokenIssuer(t, func(svc *aclsvc.Service, ctx context.Context, req *protoobject.PutRequest, cnrID cid.ID, tokens common.RequestTokens) (aclsvc.RequestInfo, error) { res, _, err := svc.PutRequestToInfo(ctx, req, req.Body.ObjectPart.(*protoobject.PutRequest_Body_Init_).Init, cnrID, acl.OpObjectPut, tokens) diff --git a/pkg/services/object/acl/v2/util.go b/pkg/services/object/acl/v2/util.go index d3b4fee24f..bab308b67a 100644 --- a/pkg/services/object/acl/v2/util.go +++ b/pkg/services/object/acl/v2/util.go @@ -17,12 +17,9 @@ func assertVerb(tok session.Object, reqVerb session.ObjectVerb) bool { return tok.AssertVerb( session.VerbObjectHead, session.VerbObjectGet, - session.VerbObjectDelete, - session.VerbObjectRange) + session.VerbObjectDelete) case session.VerbObjectSearch: return tok.AssertVerb(session.VerbObjectSearch, session.VerbObjectDelete) - case session.VerbObjectRange: - return tok.AssertVerb(session.VerbObjectRange) } } diff --git a/pkg/services/object/acl/v2/util_test.go b/pkg/services/object/acl/v2/util_test.go index 3c564d20a2..f84c04d33b 100644 --- a/pkg/services/object/acl/v2/util_test.go +++ b/pkg/services/object/acl/v2/util_test.go @@ -21,9 +21,7 @@ func TestIsVerbCompatible(t *testing.T) { session.VerbObjectHead, session.VerbObjectGet, session.VerbObjectDelete, - session.VerbObjectRange, }, - session.VerbObjectRange: {session.VerbObjectRange}, session.VerbObjectSearch: {session.VerbObjectSearch, session.VerbObjectDelete}, } @@ -31,7 +29,6 @@ func TestIsVerbCompatible(t *testing.T) { session.VerbObjectPut, session.VerbObjectDelete, session.VerbObjectHead, - session.VerbObjectRange, session.VerbObjectGet, session.VerbObjectSearch, } diff --git a/pkg/services/object/get.go b/pkg/services/object/get.go index 3ccf6090c1..d1bd860c7b 100644 --- a/pkg/services/object/get.go +++ b/pkg/services/object/get.go @@ -38,11 +38,6 @@ var ( ServerStreams: true, ClientStreams: false, } - getRangeStreamDesc = &grpc.StreamDesc{ - StreamName: "GetRange", - ServerStreams: true, - ClientStreams: false, - } ) type getStreamProgress struct { @@ -74,10 +69,6 @@ func callGet(ctx context.Context, conn *grpc.ClientConn, request any) (grpc.Clie return callServerStream(ctx, conn, protoobject.ObjectService_Get_FullMethodName, getStreamDesc, request) } -func callRange(ctx context.Context, conn *grpc.ClientConn, request any) (grpc.ClientStream, error) { - return callServerStream(ctx, conn, protoobject.ObjectService_GetRange_FullMethodName, getRangeStreamDesc, request) -} - // returns: // - nil on completed object transmission // - [object.SplitInfoError]/nil on split info response and unset/set raw flag in request @@ -1061,7 +1052,3 @@ func calculateInitGetResponseFieldLength(idLen, sigLen, hdrLen int) int { func forwardGetRequest(ctx context.Context, req any, respStream grpc.ServerStream, node clientcore.MultiAddressClient) error { return forwardServerStreamRequest(ctx, req, respStream, node, callGet) } - -func forwardRangeRequest(ctx context.Context, req any, respStream grpc.ServerStream, node clientcore.MultiAddressClient) error { - return forwardServerStreamRequest(ctx, req, respStream, node, callRange) -} diff --git a/pkg/services/object/get/ec.go b/pkg/services/object/get/ec.go index e494581f4b..3b4aaa5b62 100644 --- a/pkg/services/object/get/ec.go +++ b/pkg/services/object/get/ec.go @@ -672,17 +672,6 @@ func (s *Service) getECPartFromNode(ctx context.Context, cnr cid.ID, parent oid. return hdr, rc, nil } -// looks up for local object that carries EC part produced within cnr for parent -// object and indexed by pi, and writes its payload range into dst. Both zero -// off and ln correspond to full payload. -// -// Returns [apistatus.ErrObjectAlreadyRemoved] if the object was marked for -// removal. Returns [apistatus.ErrObjectNotFound] if the object is missing. -// Returns [apistatus.ErrObjectOutOfRange] if the range is out of payload range. -func (s *Service) copyLocalECPartRange(ctx context.Context, dst ChunkWriter, cnr cid.ID, parent oid.ID, pi iec.PartInfo, off, ln uint64) error { - return s.copyLocalECPartPayloadRange(ctx, dst, cnr, parent, pi, common.NewPayloadRange(off, ln), nil) -} - func (s *Service) copyLocalECPartPayloadRange(ctx context.Context, dst ChunkWriter, cnr cid.ID, parent oid.ID, pi iec.PartInfo, rng common.PayloadRange, headerFn func(*object.Object) error) error { hdr, pldLen, rc, err := s.localObjects.GetECPartRange(ctx, cnr, parent, pi, rng, headerFn != nil) if err != nil { diff --git a/pkg/services/object/get/exec.go b/pkg/services/object/get/exec.go index aa744d1901..7596f1bbf9 100644 --- a/pkg/services/object/get/exec.go +++ b/pkg/services/object/get/exec.go @@ -34,7 +34,7 @@ type execCtx struct { ctx context.Context - prm RangePrm + prm rangePrm payloadRange common.PayloadRange statusError @@ -72,14 +72,8 @@ type execCtx struct { getTransportFn GetTransportFunc - rangeTransportFn RangeTransportFunc - - localRangeBuffer []byte - submitLocalRangeStreamFn SubmitDataStreamFunc - payloadOnly bool recheckEACL bool - legacyRange bool // collectOnly keeps one fetched stream and skips writing/assembly. // Virtual children are reported to the caller, not assembled recursively. @@ -140,12 +134,6 @@ func withEACLRecheck(v bool) execOption { } } -func withLegacyRange(v bool) execOption { - return func(c *execCtx) { - c.legacyRange = v - } -} - func withLogger(l *zap.Logger) execOption { return func(ctx *execCtx) { ctx.log = l @@ -180,19 +168,6 @@ func withGetTransportFunc(f GetTransportFunc) execOption { } } -func withRangeTransportFunc(f RangeTransportFunc) execOption { - return func(ctx *execCtx) { - ctx.rangeTransportFn = f - } -} - -func withLocalRangeBuffer(buf []byte, submitStreamFn SubmitDataStreamFunc) execOption { - return func(ctx *execCtx) { - ctx.localRangeBuffer = buf - ctx.submitLocalRangeStreamFn = submitStreamFn - } -} - func (exec *execCtx) setLogger(l *zap.Logger) { if l.Level() != zap.DebugLevel { exec.log = l @@ -347,11 +322,10 @@ func (exec *execCtx) headOnly() bool { type childFetchCtx struct { svc *Service - prm RangePrm + prm rangePrm cnr cid.ID payloadOnly bool - legacyRange bool log *zap.Logger } @@ -364,7 +338,6 @@ func (exec *execCtx) childFetchCtx() childFetchCtx { cnr: exec.containerID(), payloadOnly: exec.payloadOnly, - legacyRange: exec.legacyRange, log: exec.log, } @@ -393,7 +366,6 @@ func (c childFetchCtx) fetchChildStream(ctx context.Context, id oid.ID, rng *obj se := c.svc.get(ctx, p.commonPrm, withPayloadRange(rng), withPayloadOnly(c.payloadOnly), - withLegacyRange(c.legacyRange), withLogger(log), withCollectOnlyResult(&res), ) @@ -438,7 +410,6 @@ func (exec *execCtx) copyChild(id oid.ID, rng *object.Range, withHdr bool) bool withPayloadRange(rng), withPayloadOnly(exec.payloadOnly), withEACLRecheck(exec.recheckEACL), - withLegacyRange(exec.legacyRange), withLogger(log), ) @@ -695,6 +666,5 @@ func (exec *execCtx) writeCollectedObject() { // parameters, so it won't be inherited in new execution contexts. func (exec *execCtx) disableForwarding() { exec.getTransportFn = nil - exec.rangeTransportFn = nil exec.headTransportFn = nil } diff --git a/pkg/services/object/get/get.go b/pkg/services/object/get/get.go index a06435a645..fe92b0116b 100644 --- a/pkg/services/object/get/get.go +++ b/pkg/services/object/get/get.go @@ -5,11 +5,8 @@ import ( "errors" "fmt" - iec "github.com/nspcc-dev/neofs-node/internal/ec" inetmap "github.com/nspcc-dev/neofs-node/internal/netmap" - "github.com/nspcc-dev/neofs-node/pkg/local_object_storage/blobstor/common" apistatus "github.com/nspcc-dev/neofs-sdk-go/client/status" - "github.com/nspcc-dev/neofs-sdk-go/netmap" "github.com/nspcc-dev/neofs-sdk-go/object" "go.uber.org/zap" ) @@ -114,71 +111,6 @@ func writeObjectHeader(dst ObjectWriter, hdr *object.Object, payloadOnly bool) e return dst.WriteHeader(hdr) } -// GetRange serves a request to get an object by address, and returns Streamer instance. -func (s *Service) GetRange(ctx context.Context, prm RangePrm) error { - pi, err := checkECPartInfoRequest(prm.common.XHeaders(), prm.container) - if err != nil { - // TODO: track https://github.com/nspcc-dev/neofs-api/issues/269. - return fmt.Errorf("invalid request: %w", err) - } - - if pi.RuleIndex >= 0 { - // TODO: deny if node is not in the container? - - if prm.localBuffer != nil { - stream, err := s.localObjects.ReadECPartRange(ctx, prm.addr.Container(), prm.addr.Object(), pi, prm.rng.GetOffset(), prm.rng.GetLength(), prm.localBuffer, nil) - if err == nil { - prm.submitLocalStreamFn(stream) - } - return err - } - - return s.copyLocalECPartRange(ctx, prm.objWriter, prm.addr.Container(), prm.addr.Object(), pi, prm.rng.GetOffset(), prm.rng.GetLength()) - } - - if prm.common.LocalOnly() && - len(prm.container.PlacementPolicy().ECRules()) == 0 && // EC breaks TTL requirements currently. - len(prm.container.PlacementPolicy().Replicas()) != 0 { - // It handles locality internally. - bufOpt := withLocalRangeBuffer(prm.localBuffer, prm.submitLocalStreamFn) - return s.get(ctx, prm.commonPrm, withPayloadRange(prm.rng), withPayloadOnly(true), withLegacyRange(true), bufOpt).err - } - - nodeLists, repRules, ecRules, err := s.neoFSNet.GetNodesForObject(prm.addr) - if err != nil { - return fmt.Errorf("get nodes for object: %w", err) - } - - if prm.forwardRequestFn != nil && !inetmap.NodeSetsContainPublicKeyFunc(nodeLists, s.neoFSNet.IsLocalNodePublicKey) { - return s.forwardRequest(ctx, repRules, ecRules, nodeLists, "RANGE", prm.forwardRequestFn) - } - - return s.getRange(ctx, prm, nodeLists, repRules, ecRules) -} - -func (s *Service) getRange(ctx context.Context, prm RangePrm, nodeLists [][]netmap.NodeInfo, repRules []uint, ecRules []iec.Rule) error { - if len(repRules) > 0 { // REP format does not require encoding - bufOpt := withLocalRangeBuffer(prm.localBuffer, prm.submitLocalStreamFn) - transportOpt := withRangeTransportFunc(prm.transportFn) - err := s.get(ctx, prm.commonPrm, withPreSortedContainerNodes(nodeLists[:len(repRules)], repRules), withPayloadRange(prm.rng), withPayloadOnly(true), withLegacyRange(true), bufOpt, transportOpt).err - if len(ecRules) == 0 || !errors.Is(err, apistatus.ErrObjectNotFound) { - return err - } - } - - ecNodeLists := nodeLists[len(repRules):] - - if prm.raw { - repRules = make([]uint, len(ecRules)) - for i := range ecRules { - repRules[i] = uint(ecRules[i].DataPartNum + ecRules[i].ParityPartNum) - } - return s.get(ctx, prm.commonPrm, withPreSortedContainerNodes(ecNodeLists, repRules), withPayloadRange(prm.rng), withPayloadOnly(true), withLegacyRange(true)).err - } - - return s.copyECObjectRange(ctx, prm.objWriter, prm.addr.Container(), prm.addr.Object(), ecRules, ecNodeLists, common.NewPayloadRange(prm.rng.GetOffset(), prm.rng.GetLength()), nil) -} - // Head reads object header from container. // // Returns ErrNotFound if the header was not received for the call. @@ -251,7 +183,7 @@ func (s *Service) get(ctx context.Context, prm commonPrm, opts ...execOption) st exec := &execCtx{ svc: s, ctx: ctx, - prm: RangePrm{ + prm: rangePrm{ commonPrm: prm, }, infoSplit: object.NewSplitInfo(), diff --git a/pkg/services/object/get/get_test.go b/pkg/services/object/get/get_test.go index 714a8a2870..f2ec29b14f 100644 --- a/pkg/services/object/get/get_test.go +++ b/pkg/services/object/get/get_test.go @@ -225,7 +225,7 @@ func (s *testStorage) get(exec *execCtx) (*object.Object, io.ReadCloser, error) func (s *testStorage) Head(_ context.Context, addr oid.Address, _ bool) (*object.Object, error) { hdr, _, err := s.get(&execCtx{ - prm: RangePrm{ + prm: rangePrm{ commonPrm: commonPrm{ addr: addr, }, @@ -304,21 +304,6 @@ func TestGetLocalOnly(t *testing.T) { return p } - newRngPrm := func(raw bool, w ChunkWriter, off, ln uint64) RangePrm { - p := RangePrm{} - p.SetChunkWriter(w) - p.WithRawFlag(raw) - p.common = new(util.CommonPrm).WithLocalOnly(true) - - r := object.NewRange() - r.SetOffset(off) - r.SetLength(ln) - - p.SetRange(r) - - return p - } - newHeadPrm := func(raw bool, w ObjectWriter) HeadPrm { p := HeadPrm{} p.SetHeaderWriter(w) @@ -353,15 +338,6 @@ func TestGetLocalOnly(t *testing.T) { require.Equal(t, obj, w.Object()) - w = NewSimpleObjectWriter() - - rngPrm := newRngPrm(false, w, payloadSz/3, payloadSz/3) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.NoError(t, err) - require.Equal(t, payload[payloadSz/3:2*payloadSz/3], w.Object().Payload()) - w = NewSimpleObjectWriter() headPrm := newHeadPrm(false, w) headPrm.WithAddress(addr) @@ -385,12 +361,6 @@ func TestGetLocalOnly(t *testing.T) { require.ErrorAs(t, err, new(apistatus.ObjectAlreadyRemoved)) - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectAlreadyRemoved)) - headPrm := newHeadPrm(false, nil) headPrm.WithAddress(addr) @@ -410,13 +380,6 @@ func TestGetLocalOnly(t *testing.T) { require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - headPrm := newHeadPrm(false, nil) headPrm.WithAddress(addr) @@ -443,13 +406,6 @@ func TestGetLocalOnly(t *testing.T) { require.Equal(t, si, errSplit.SplitInfo()) - rngPrm := newRngPrm(true, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.Get(ctx, p) - - require.True(t, errors.As(err, &errSplit)) - headPrm := newHeadPrm(true, nil) headPrm.WithAddress(addr) @@ -715,21 +671,6 @@ func TestGetRemoteSmall(t *testing.T) { return p } - newRngPrm := func(raw bool, w ChunkWriter, off, ln uint64) RangePrm { - p := RangePrm{} - p.SetChunkWriter(w) - p.WithRawFlag(raw) - p.common = new(util.CommonPrm).WithLocalOnly(false) - - r := object.NewRange() - r.SetOffset(off) - r.SetLength(ln) - - p.SetRange(r) - - return p - } - newHeadPrm := func(raw bool, w ObjectWriter) HeadPrm { p := HeadPrm{} p.SetHeaderWriter(w) @@ -783,14 +724,6 @@ func TestGetRemoteSmall(t *testing.T) { require.NoError(t, err) require.Equal(t, obj, w.Object()) - w = NewSimpleObjectWriter() - rngPrm := newRngPrm(false, w, payloadSz/3, payloadSz/3) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.NoError(t, err) - require.Equal(t, payload[payloadSz/3:2*payloadSz/3], w.Object().Payload()) - w = NewSimpleObjectWriter() headPrm := newHeadPrm(false, w) headPrm.WithAddress(addr) @@ -829,12 +762,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(*apistatus.ObjectAlreadyRemoved)) - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(*apistatus.ObjectAlreadyRemoved)) - headPrm := newHeadPrm(false, nil) headPrm.WithAddress(addr) @@ -871,12 +798,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - headPrm := newHeadPrm(false, nil) headPrm.WithAddress(addr) @@ -940,12 +861,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) }) t.Run("get chain element failure", func(t *testing.T) { @@ -1014,12 +929,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - - rngPrm := newRngPrm(false, NewSimpleObjectWriter(), 0, 1) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) }) t.Run("OK", func(t *testing.T) { @@ -1095,22 +1004,13 @@ func TestGetRemoteSmall(t *testing.T) { require.Equal(t, srcObj, w.Object()) w = NewSimpleObjectWriter() - payloadSz := srcObj.PayloadSize() + p = newPrm(false, w) + p.WithAddress(addr) + payloadSz := srcObj.PayloadSize() off := payloadSz / 3 ln := payloadSz / 3 - rngPrm := newRngPrm(false, w, off, ln) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.NoError(t, err) - require.Equal(t, payload[off:off+ln], w.Object().Payload()) - - w = NewSimpleObjectWriter() - p = newPrm(false, w) - p.WithAddress(addr) - r := object.NewRange() r.SetOffset(off) r.SetLength(ln) @@ -1167,12 +1067,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - - rngPrm := newRngPrm(false, nil, 0, 0) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) }) t.Run("get chain element failure", func(t *testing.T) { @@ -1234,12 +1128,6 @@ func TestGetRemoteSmall(t *testing.T) { err := svc.Get(ctx, p) require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) - - rngPrm := newRngPrm(false, nil, 0, 1) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.ErrorAs(t, err, new(apistatus.ObjectNotFound)) }) t.Run("OK", func(t *testing.T) { @@ -1307,22 +1195,13 @@ func TestGetRemoteSmall(t *testing.T) { require.Equal(t, srcObj, w.Object()) w = NewSimpleObjectWriter() - payloadSz := srcObj.PayloadSize() + p = newPrm(false, w) + p.WithAddress(addr) + payloadSz := srcObj.PayloadSize() off := payloadSz / 3 ln := payloadSz / 3 - rngPrm := newRngPrm(false, w, off, ln) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.NoError(t, err) - require.Equal(t, payload[off:off+ln], w.Object().Payload()) - - w = NewSimpleObjectWriter() - p = newPrm(false, w) - p.WithAddress(addr) - r := object.NewRange() r.SetOffset(off) r.SetLength(ln) @@ -1332,17 +1211,6 @@ func TestGetRemoteSmall(t *testing.T) { require.NoError(t, err) require.Equal(t, srcObj.CutPayload(), w.Object().CutPayload()) require.Equal(t, payload[off:off+ln], w.Object().Payload()) - - w = NewSimpleObjectWriter() - off = payloadSz - 2 - ln = 1 - - rngPrm = newRngPrm(false, w, off, ln) - rngPrm.WithAddress(addr) - - err = svc.GetRange(ctx, rngPrm) - require.NoError(t, err) - require.Equal(t, payload[off:off+ln], w.Object().Payload()) }) }) }) @@ -1364,7 +1232,7 @@ func TestWriteCollectedHeaderPayloadOnlyDoesNotStartResponse(t *testing.T) { t.Run("valid header", func(t *testing.T) { exec := &execCtx{ - prm: RangePrm{ + prm: rangePrm{ commonPrm: commonPrm{ objWriter: &trackingWriter{}, }, @@ -1387,7 +1255,7 @@ func TestWriteCollectedHeaderPayloadOnlyDoesNotStartResponse(t *testing.T) { t.Run("missing header", func(t *testing.T) { exec := &execCtx{ - prm: RangePrm{ + prm: rangePrm{ commonPrm: commonPrm{ objWriter: &trackingWriter{}, }, @@ -1404,59 +1272,6 @@ func TestWriteCollectedHeaderPayloadOnlyDoesNotStartResponse(t *testing.T) { }) } -func TestFallbackRangeReader(t *testing.T) { - t.Run("fallback get exec clears range and flags", func(t *testing.T) { - rng := object.NewRange() - rng.SetOffset(10) - rng.SetLength(20) - - exec := execCtx{ - payloadRange: blobcommon.NewPayloadRange(rng.GetOffset(), rng.GetLength()), - payloadOnly: true, - legacyRange: true, - } - - fallback := exec.fallbackGetExec() - require.Nil(t, fallback.ctxRange()) - require.False(t, fallback.payloadOnly) - require.False(t, fallback.legacyRange) - }) - - t.Run("partial chunk is returned before fallback", func(t *testing.T) { - buf := make([]byte, 16) - fr := &fallbackRangeReader{ - ReadCloser: &partialErrorReader{ - data: []byte("payload"), - err: apistatus.ErrObjectAccessDenied, - }, - } - - n, err := fr.Read(buf) - require.NoError(t, err) - require.Equal(t, len("payload"), n) - require.Equal(t, []byte("payload"), buf[:n]) - require.EqualValues(t, len("payload"), fr.delivered) - require.True(t, fr.fallbackPending) - require.False(t, fr.fallbackDone) - }) - - t.Run("fallback resumes after already delivered bytes", func(t *testing.T) { - rng := object.NewRange() - rng.SetOffset(10) - rng.SetLength(20) - - fr := &fallbackRangeReader{ - rng: rng, - delivered: 7, - } - - from, to, err := fr.fallbackBounds(100) - require.NoError(t, err) - require.EqualValues(t, 17, from) - require.EqualValues(t, 30, to) - }) -} - type failingReader struct { data []byte pos int @@ -1494,24 +1309,6 @@ func (r *errorReader) Close() error { return nil } -type partialErrorReader struct { - data []byte - err error - read bool -} - -func (r *partialErrorReader) Read(p []byte) (int, error) { - if r.read { - return 0, io.EOF - } - r.read = true - return copy(p, r.data), r.err -} - -func (r *partialErrorReader) Close() error { - return nil -} - type trackingWriter struct { writeHeaderCount atomic.Int32 writeChunkCount atomic.Int32 diff --git a/pkg/services/object/get/prm.go b/pkg/services/object/get/prm.go index 39c4e417ad..34f4735a8e 100644 --- a/pkg/services/object/get/prm.go +++ b/pkg/services/object/get/prm.go @@ -41,16 +41,11 @@ type Prm struct { interceptHeaderBinaryFn func([]byte) error } -// RangePrm groups parameters of GetRange service call. -type RangePrm struct { +// rangePrm groups common parameters with range definition. +type rangePrm struct { commonPrm rng *object.Range - - localBuffer []byte - submitLocalStreamFn SubmitDataStreamFunc - - transportFn RangeTransportFunc } type RequestForwarder func(context.Context, clientcore.MultiAddressClient) (*object.Object, error) @@ -67,10 +62,6 @@ type SubmitHeadResponseFunc = func(mem.BufferSlice, iprotobuf.BuffersSlice) // through passed connection. type GetTransportFunc func(context.Context, clientcore.MultiAddressClient) error -// RangeTransportFunc continues to serve current RANGE request from remote node -// through passed connection. -type RangeTransportFunc func(context.Context, clientcore.MultiAddressClient) error - // HeadPrm groups parameters of Head service call. type HeadPrm struct { commonPrm @@ -169,15 +160,8 @@ func (p *Prm) RequireEACLRecheck() { p.recheckEACL = true } -// SetChunkWriter sets target component to write the object payload range. -func (p *RangePrm) SetChunkWriter(w ChunkWriter) { - p.objWriter = &partWriter{ - chunkWriter: w, - } -} - // SetRange sets range of the requested payload data. -func (p *RangePrm) SetRange(rng *object.Range) { +func (p *rangePrm) SetRange(rng *object.Range) { p.rng = rng } @@ -306,27 +290,6 @@ func (p *Prm) SetTransportFunc(f GetTransportFunc) { p.transportFn = f } -// WithBuffer specifies a buffer to use for header reading and a callback for -// payload range stream. If passed, the stream must be finally closed by the -// caller. -func (p *RangePrm) WithBuffer(buffer []byte, submitStreamFn SubmitDataStreamFunc) { - p.localBuffer = buffer - p.submitLocalStreamFn = submitStreamFn -} - -// SetTransportFunc specifies request transport callback to use for streaming -// responses from remote node by in-container server. -// -// The f should return: -// - nil on completed object transmission -// - [object.SplitInfoError] on OK with corresponding body field -// - [apistatus.ErrObjectNotFound] on 404 status -// - nil on other API statuses -// - any other transport/protocol error otherwise -func (p *RangePrm) SetTransportFunc(f RangeTransportFunc) { - p.transportFn = f -} - // WithECTransport specifies transport layer to for EC handling. func (p *Prm) WithECTransport(transport GetECRequestTransport) { p.ecTransport = transport diff --git a/pkg/services/object/get/util.go b/pkg/services/object/get/util.go index 0ed1c1c7ce..3c74b080af 100644 --- a/pkg/services/object/get/util.go +++ b/pkg/services/object/get/util.go @@ -14,7 +14,6 @@ import ( "github.com/nspcc-dev/neofs-node/pkg/services/object/internal" "github.com/nspcc-dev/neofs-sdk-go/bearer" "github.com/nspcc-dev/neofs-sdk-go/client" - apistatus "github.com/nspcc-dev/neofs-sdk-go/client/status" cid "github.com/nspcc-dev/neofs-sdk-go/container/id" "github.com/nspcc-dev/neofs-sdk-go/netmap" "github.com/nspcc-dev/neofs-sdk-go/object" @@ -59,12 +58,9 @@ type objectReadAuthPrm interface { WithBearerToken(bearer.Token) } -func applyObjectReadAuth(exec *execCtx, addr oid.Address, legacyRange bool, opts objectReadAuthPrm) { +func applyObjectReadAuth(exec *execCtx, addr oid.Address, opts objectReadAuthPrm) { if stV2 := exec.prm.common.SessionTokenV2(); stV2 != nil { verb := sessionv2.VerbObjectGet - if legacyRange { - verb = sessionv2.VerbObjectRange - } if stV2.AssertVerb(verb, addr.Container()) { opts.WithinSessionV2(*stV2) } @@ -96,123 +92,6 @@ type partWriter struct { chunkWriter ChunkWriter } -// fallbackRangeReader wraps a range reader obtained via ObjectRangeInit and -// falls back to a full GET in case apistatus.ErrObjectAccessDenied is -// returned while reading. -type fallbackRangeReader struct { - io.ReadCloser - exec *execCtx - client *clientWrapper - key *ecdsa.PrivateKey - rng *object.Range - - delivered uint64 - fallbackPending bool - fallbackDone bool -} - -func (exec execCtx) fallbackGetExec() execCtx { - exec.payloadRange = common.PayloadRange{} - exec.payloadOnly = false - exec.legacyRange = false - - return exec -} - -func newFallbackRangeReader(exec *execCtx, c *clientWrapper, key *ecdsa.PrivateKey, rng *object.Range, rdr io.ReadCloser) io.ReadCloser { - return &fallbackRangeReader{ - ReadCloser: rdr, - exec: exec, - client: c, - key: key, - rng: rng, - } -} - -func (f *fallbackRangeReader) Read(p []byte) (int, error) { - if f.fallbackPending && !f.fallbackDone { - return f.fallbackRead(p) - } - - n, err := f.ReadCloser.Read(p) - f.delivered += uint64(n) - if err == nil || !errors.Is(err, apistatus.ErrObjectAccessDenied) || f.fallbackDone { - return n, err - } - if n > 0 { - f.fallbackPending = true - return n, nil - } - - return f.fallbackRead(p) -} - -func (f *fallbackRangeReader) fallbackBounds(payloadSize uint64) (from, to uint64, err error) { - base := f.rng.GetOffset() - from = base + f.delivered - if from < base { - return 0, 0, apistatus.ErrObjectOutOfRange - } - - if ln := f.rng.GetLength(); ln != 0 { - to = base + ln - if to < base || to < from { - return 0, 0, apistatus.ErrObjectOutOfRange - } - } else { - to = payloadSize - } - - if payloadSize < from || payloadSize < to { - return 0, 0, apistatus.ErrObjectOutOfRange - } - - return from, to, nil -} - -func (f *fallbackRangeReader) fallbackRead(p []byte) (int, error) { - // TODO: drop fallback once legacy RANGE is aligned with GET access semantics, see #3547. - f.exec.log.Debug("range read access denied, falling back to full GET") - f.fallbackPending = false - f.fallbackDone = true - - oldRdr := f.ReadCloser - if oldRdr != nil { - defer func() { _ = oldRdr.Close() }() - } - - fallbackExec := f.exec.fallbackGetExec() - hdr, rdr, getErr := f.client.get(&fallbackExec, f.key) - if getErr != nil { - return 0, fmt.Errorf("fallback GET after access denial failed: %w", getErr) - } - - from, to, err := f.fallbackBounds(hdr.PayloadSize()) - if err != nil { - _ = rdr.Close() - return 0, err - } - - if from > 0 { - _, err = io.CopyN(io.Discard, rdr, int64(from)) - if err != nil { - _ = rdr.Close() - return 0, fmt.Errorf("discard %d bytes in stream: %w", from, err) - } - } - - f.ReadCloser = struct { - io.Reader - io.Closer - }{ - Reader: io.LimitReader(rdr, int64(to-from)), - Closer: rdr, - } - - // attempt to read again immediately to fill p. - return f.Read(p) -} - func NewSimpleObjectWriter() *SimpleObjectWriter { return &SimpleObjectWriter{ obj: new(object.Object), @@ -264,10 +143,6 @@ func (c *clientWrapper) getObject(exec *execCtx) (*object.Object, io.ReadCloser, return nil, nil, exec.getTransportFn(exec.ctx, c.client) } - if exec.rangeTransportFn != nil { - return nil, nil, exec.rangeTransportFn(exec.ctx, c.client) - } - key, err := exec.key() if err != nil { return nil, nil, err @@ -285,39 +160,11 @@ func (c *clientWrapper) getObject(exec *execCtx) (*object.Object, io.ReadCloser, // we don't specify payload writer because we accumulate // the object locally (even huge). if exec.hasPayloadRange() { - rng := exec.ctxRange() addr := exec.address() id := addr.Object() - if exec.legacyRange { - ln := rng.GetLength() - - var opts client.PrmObjectRange - if exec.prm.common.TTL() < 2 { - opts.MarkLocal() - } - applyObjectReadAuth(exec, addr, true, &opts) - opts.WithXHeaders(exec.prm.common.XHeaders()...) - if exec.isRaw() { - opts.MarkRaw() - } - - rdr, err := c.client.ObjectRangeInit(exec.context(), addr.Container(), id, rng.GetOffset(), ln, user.NewAutoIDSigner(*key), opts) - if err != nil { - return nil, nil, fmt.Errorf("init payload reading: %w", err) - } - - hdr, err := c.head(exec, key) - if err != nil { - _ = rdr.Close() - return nil, nil, err - } - - return hdr, newFallbackRangeReader(exec, c, key, rng, rdr), nil - } - opts := objectGetOptions(exec) - applyObjectReadAuth(exec, addr, false, &opts) + applyObjectReadAuth(exec, addr, &opts) first, second := exec.payloadRange.First, exec.payloadRange.Second switch exec.payloadRange.Mode { case common.PayloadRangeModeNone: @@ -399,15 +246,6 @@ func (e *storageEngineWrapper) get(exec *execCtx) (*object.Object, io.ReadCloser } if exec.hasPayloadRange() { - if exec.localRangeBuffer != nil { - rng := exec.ctxRange() - r, err := e.engine.ReadPayloadRange(ctx, exec.address(), rng.GetOffset(), rng.GetLength(), exec.localRangeBuffer) - if err == nil { - exec.submitLocalRangeStreamFn(r) - } - return nil, nil, err - } - hdr, stream, err := e.engine.GetRangeStream(ctx, exec.address(), exec.payloadRange, true) if err != nil { return nil, stream, err diff --git a/pkg/services/object/proto.go b/pkg/services/object/proto.go index 2045e80db8..99d9fa87c4 100644 --- a/pkg/services/object/proto.go +++ b/pkg/services/object/proto.go @@ -226,10 +226,6 @@ func shiftPayloadChunkInGetResponseBuffer(respBuf []byte, off, ln int) iprotobuf return shiftPayloadChunkInResponseBuffer(respBuf, iprotobuf.TagBytes2, off, ln) } -func shiftPayloadChunkInRangeResponseBuffer(respBuf []byte, off, ln int) iprotobuf.FieldBounds { - return shiftPayloadChunkInResponseBuffer(respBuf, iprotobuf.TagBytes1, off, ln) -} - func shiftPayloadChunkInResponseBuffer(respBuf []byte, chunkFldTag byte, off, ln int) iprotobuf.FieldBounds { bodyFldPrefixLen := 1 + protowire.SizeVarint(uint64(ln)) diff --git a/pkg/services/object/range.go b/pkg/services/object/range.go deleted file mode 100644 index 278bfb0ca7..0000000000 --- a/pkg/services/object/range.go +++ /dev/null @@ -1,179 +0,0 @@ -package object - -import ( - "context" - "errors" - "fmt" - "io" - - apistatus "github.com/nspcc-dev/neofs-sdk-go/client/status" - 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" - protostatus "github.com/nspcc-dev/neofs-sdk-go/proto/status" - "google.golang.org/grpc" - "google.golang.org/grpc/mem" - "google.golang.org/protobuf/encoding/protowire" -) - -type rangeStreamProgress struct { - empty bool - readPayload int -} - -// returns: -// - nil on completed payload transmission -// - [object.SplitInfoError]/nil on split info response and unset/set raw flag in request -// - [apistatus.ErrObjectNotFound] on 404 status -// - nil on other API statuses -// - any other transport/protocol error otherwise -func (s *rangeStream) continueWithConn(ctx context.Context, conn *grpc.ClientConn) error { - stream, err := conn.NewStream(ctx, &protoobject.ObjectService_ServiceDesc.Streams[3], protoobject.ObjectService_GetRange_FullMethodName, - grpc.StaticMethod(), - grpc.ForceCodecV2(iprotobuf.BufferedCodec{}), - ) - if err != nil { - return fmt.Errorf("stream opening failed: %w", err) - } - if err = stream.SendMsg(s.req); err != nil { - return fmt.Errorf("send request: %w", err) - } - if err = stream.CloseSend(); err != nil { - return fmt.Errorf("close send: %w", err) - } - - var prog rangeStreamProgress - prog.empty = true - - for { - var respBuf mem.BufferSlice - if err = stream.RecvMsg(&respBuf); err != nil { - if errors.Is(err, io.EOF) { - if prog.empty { - return io.ErrUnexpectedEOF - } - return nil - } - return fmt.Errorf("reading the response failed: %w", err) - } - - prog.empty = false - - fin, sent, err := s.handleResponse(&prog, respBuf) - if !sent { - respBuf.Free() - } - if err != nil { - return fmt.Errorf("handle next stream message: %w", err) - } - if fin { - return nil - } - } -} - -func (s *rangeStream) handleResponse(streamProg *rangeStreamProgress, respBuf mem.BufferSlice) (bool, bool, error) { - var code uint32 - var body iprotobuf.BuffersSlice - - var opts protoscan.ScanMessageOptions - opts.InterceptNested = func(num protowire.Number, buffers iprotobuf.BuffersSlice) error { - switch num { - default: - return protoscan.ErrContinue - case iprotobuf.FieldResponseBody: - body = buffers - return nil - case iprotobuf.FieldResponseMetaHeader: - var err error - code, err = getStatusCodeFromResponseMetaHeader(buffers) - if err != nil { - return fmt.Errorf("handle meta header: %w", err) - } - return nil - } - } - - err := protoscan.ScanMessage(iprotobuf.NewBuffersSlice(respBuf), protoscan.ResponseScheme, opts) - if err != nil { - return false, false, err - } - - if code == protostatus.ObjectNotFound { - return false, false, apistatus.ErrObjectNotFound - } - - if code != protostatus.OK { - return true, true, s.base.SendMsg(respBuf) - } - - sent, err := s.handleResponseBody(streamProg, respBuf, body) - if err != nil { - return false, sent, fmt.Errorf("handle body: %w", err) - } - - return false, sent, nil -} - -func (s *rangeStream) handleResponseBody(streamProg *rangeStreamProgress, respBuf mem.BufferSlice, buffers iprotobuf.BuffersSlice) (bool, error) { - var oneofNum protowire.Number - var oneofFld iprotobuf.BuffersSlice - - var opts protoscan.ScanMessageOptions - opts.InterceptBytes = func(num protowire.Number, buffers iprotobuf.BuffersSlice) error { - if num == protoobject.FieldRangeResponseBodyChunk { - oneofNum, oneofFld = num, buffers - } - return nil - } - opts.InterceptNested = func(num protowire.Number, buffers iprotobuf.BuffersSlice) error { - switch num { - default: - return protoscan.ErrContinue - case protoobject.FieldRangeResponseBodySplitInfo: - oneofNum, oneofFld = num, buffers - return nil - } - } - - err := protoscan.ScanMessage(buffers, protoscan.ObjectGetRangeResponseBodyScheme, opts) - if err != nil { - return false, err - } - - switch oneofNum { - default: - return false, errors.New("none of the supported oneof fields are specified") - case protoobject.FieldRangeResponseBodyChunk: - return s.handleChunkResponse(streamProg, respBuf, oneofFld) - case protoobject.FieldRangeResponseBodySplitInfo: - return s.handleSplitInfo(respBuf, oneofFld) - } -} - -func (s *rangeStream) handleChunkResponse(streamProg *rangeStreamProgress, respBuf mem.BufferSlice, chunkBuffers iprotobuf.BuffersSlice) (bool, error) { - chunkLen := chunkBuffers.Len() - - from, to := chunkBoundsToSend(s.respondedPayload, streamProg.readPayload, chunkLen) - if from == to { - streamProg.readPayload += chunkLen - return false, nil - } - - _, ok := chunkBuffers.MoveNext(from) - if !ok { - return false, fmt.Errorf("seek chunk left bound in response buffers: %w", io.ErrUnexpectedEOF) - } - - chunkBuffers, ok = chunkBuffers.MoveNext(to - from) - if !ok { - return false, fmt.Errorf("seek chunk right bound in response buffers: %w", io.ErrUnexpectedEOF) - } - - return s.srv.sendChunkResponse(s.base, respBuf, chunkBuffers, to-from, chunkLen, - s.signResponse, iprotobuf.TagBytes1, &streamProg.readPayload, &s.respondedPayload, shiftPayloadChunkInRangeResponseBuffer) -} - -func (s *rangeStream) handleSplitInfo(respBuf mem.BufferSlice, buffers iprotobuf.BuffersSlice) (bool, error) { - return handleSplitInfoAndRespond(s.req.GetBody().GetRaw(), s.base, respBuf, buffers) -} diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index 78292d6a52..4421f23be5 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -70,7 +70,6 @@ type Handlers interface { Put(context.Context) (*putsvc.Streamer, error) Head(context.Context, getsvc.HeadPrm) error Delete(context.Context, deletesvc.Prm) error - GetRange(context.Context, getsvc.RangePrm) error } // Various NeoFS protocol status codes. @@ -171,7 +170,6 @@ type ACLInfoExtractor interface { DeleteRequestToInfo(context.Context, *protoobject.DeleteRequest, cid.ID, common.RequestTokens) (aclsvc.RequestInfo, error) HeadRequestToInfo(context.Context, *protoobject.HeadRequest, cid.ID, common.RequestTokens) (aclsvc.RequestInfo, error) GetRequestToInfo(context.Context, *protoobject.GetRequest, cid.ID, common.RequestTokens) (aclsvc.RequestInfo, error) - RangeRequestToInfo(context.Context, *protoobject.GetRangeRequest, cid.ID, common.RequestTokens) (aclsvc.RequestInfo, error) SearchV2RequestToInfo(context.Context, *protoobject.SearchV2Request, cid.ID, common.RequestTokens) (aclsvc.RequestInfo, error) VerifySessionTokenMessage(*protosession.SessionTokenV2, sessionv2.Verb, cid.ID) (sessionv2.Token, error) VerifySessionV1TokenMessage(*protosession.SessionToken, session.ObjectVerb, cid.ID, oid.ID) (session.Object, error) @@ -1506,147 +1504,8 @@ type getProxyContext struct { resolveRange func(uint64) (uint64, uint64, error) } -func (s *Server) sendRangeResponse(stream protoobject.ObjectService_GetRangeServer, resp *protoobject.GetRangeResponse, req *protoobject.GetRangeRequest) error { - resp.VerifyHeader = util.SignResponseIfNeeded(&s.signer, resp, req) - return stream.Send(resp) -} - -func (s *Server) sendStatusRangeResponse(stream protoobject.ObjectService_GetRangeServer, err error, req *protoobject.GetRangeRequest) error { - if splitErr, ok := errors.AsType[*object.SplitInfoError](err); ok { - return s.sendRangeResponse(stream, &protoobject.GetRangeResponse{ - Body: &protoobject.GetRangeResponse_Body{ - RangePart: &protoobject.GetRangeResponse_Body_SplitInfo{ - SplitInfo: splitErr.SplitInfo().ProtoMessage(), - }, - }, - }, req) - } - return s.sendRangeResponse(stream, &protoobject.GetRangeResponse{ - MetaHeader: s.makeResponseMetaHeader(util.ToStatus(err), req.MetaHeader), - }, req) -} - -type rangeStream struct { - base protoobject.ObjectService_GetRangeServer - srv *Server - req *protoobject.GetRangeRequest - - respondedPayload int - - signResponse bool -} - -func (s *rangeStream) WriteChunk(chunk []byte) error { - for buf := bytes.NewBuffer(chunk); buf.Len() > 0; { - newResp := &protoobject.GetRangeResponse{ - Body: &protoobject.GetRangeResponse_Body{ - RangePart: &protoobject.GetRangeResponse_Body_Chunk{ - Chunk: buf.Next(maxRespDataChunkSize), - }, - }, - } - if err := s.srv.sendRangeResponse(s.base, newResp, s.req); err != nil { - return err - } - } - return nil -} - -func (s *Server) GetRange(req *protoobject.GetRangeRequest, gStream protoobject.ObjectService_GetRangeServer) error { - ctx := gStream.Context() - var ( - err error - t = time.Now() - ) - defer func() { s.pushOpExecResult(stat.MethodObjectRange, err, t) }() - if err = icrypto.VerifyRequestSignaturesN3(ctx, req, s.fsChain); err != nil { - return s.sendStatusRangeResponse(gStream, err, req) - } - - if s.fsChain.LocalNodeUnderMaintenance() { - return s.sendStatusRangeResponse(gStream, apistatus.ErrNodeUnderMaintenance, req) - } - - body := req.Body - if body == nil { - err = newBadRequestError(missingRequestBodyMessage) // defer - return s.sendStatusRangeResponse(gStream, err, req) - } - - cnrID, objID, err := fetchRequiredObjectAddress(body.Address) - if err != nil { - err = newBadRequestError(invalidRequestBodyMessage + ": " + err.Error()) // defer - return s.sendStatusRangeResponse(gStream, err, req) - } - - reqMD, err := s.handleRequestMetaHeader(req.MetaHeader, sessionv2.VerbObjectRange, session.VerbObjectRange, cnrID, objID) - if err != nil { - return s.sendStatusRangeResponse(gStream, err, req) - } - - reqInfo, err := s.reqInfoProc.RangeRequestToInfo(ctx, req, cnrID, reqMD.tokens) - if err != nil { - if !errors.Is(err, apistatus.Error) { - err = newBadRequestError(err.Error()) // defer - } - return s.sendStatusRangeResponse(gStream, err, req) - } - if !s.aclChecker.CheckBasicACL(reqInfo) { - err = basicACLErr(reqInfo) // needed for defer - return s.sendStatusRangeResponse(gStream, err, req) - } - 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.sendStatusRangeResponse(gStream, err, req) - } - - needSignResponse := needSignGetResponse(req) - - p, err := convertRangePrm(s.signer, reqInfo.Container, req, &rangeStream{ - base: gStream, - srv: s, - req: req, - signResponse: needSignResponse, - }, cnrID, objID, reqMD) - if err != nil { - if !errors.Is(err, apistatus.Error) { - err = newBadRequestError(err.Error()) // defer - } - return s.sendStatusRangeResponse(gStream, err, req) - } - - p.SetForwardRequestFunc(func(ctx context.Context, node clientcore.MultiAddressClient) error { - return forwardRangeRequest(ctx, req, gStream, node) - }) - - var stream io.ReadCloser - defer func() { - if stream != nil { - stream.Close() - } - }() - - hdrRespBuf, hdrBuf := getBufferForHeadResponse() - defer hdrRespBuf.Free() - - p.WithBuffer(hdrBuf, func(s io.ReadCloser) { stream = s }) - - err = s.handlers.GetRange(ctx, p) - if err != nil { - return s.sendStatusRangeResponse(gStream, err, req) - } - - if stream == nil { - return nil - } - - err = s.copyRangeStream(gStream, stream, needSignResponse, shiftPayloadChunkInRangeResponseBuffer) - if err != nil { - return s.sendStatusRangeResponse(gStream, err, req) - } - - return nil +func (s *Server) GetRange(_ *protoobject.GetRangeRequest, _ protoobject.ObjectService_GetRangeServer) error { + return grpcstatus.Error(grpccodes.Unimplemented, "no longer supported, use Get with range options") } func (s *Server) copyRangeStream(gStream grpc.ServerStream, stream io.Reader, needSignResp bool, shiftFn func([]byte, int, int) iprotobuf.FieldBounds) error { @@ -1694,63 +1553,6 @@ func (s *Server) copyRangeStream(gStream grpc.ServerStream, stream io.Reader, ne } } -// 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 convertRangePrm(signer ecdsa.PrivateKey, cnr container.Container, req *protoobject.GetRangeRequest, stream *rangeStream, cnrID cid.ID, objID oid.ID, reqMD requestMetadata) (getsvc.RangePrm, error) { - body := req.GetBody() - - rln := body.Range.GetLength() - if rln == 0 { // includes nil range - if body.Range.Offset != 0 { - return getsvc.RangePrm{}, errors.New("zero range length") - } // else whole payload - } else if body.Range.Offset+rln <= body.Range.Offset { - return getsvc.RangePrm{}, errors.New("range overflow") - } - - cp := objutil.CommonPrmFromRequest(reqMD.ttl, reqMD.xHeaders, reqMD.tokens) - - var p getsvc.RangePrm - p.SetCommonParameters(cp) - p.WithAddress(oid.NewAddress(cnrID, objID)) - p.WithContainer(cnr) - p.WithRawFlag(body.Raw) - p.SetChunkWriter(stream) - var rng object.Range - rng.SetOffset(body.Range.Offset) - rng.SetLength(rln) - p.SetRange(&rng) - if cp.LocalOnly() { - return p, nil - } - - var onceResign sync.Once - meta := req.GetMetaHeader() - if meta == nil { - return getsvc.RangePrm{}, errors.New("missing meta header") - } - p.SetTransportFunc(func(ctx context.Context, c clientcore.MultiAddressClient) error { - var err error - onceResign.Do(func() { - req = &protoobject.GetRangeRequest{ - Body: req.Body, - MetaHeader: &protosession.RequestMetaHeader{ - Version: c.APIVersion(), - Ttl: 1, - }, - } - req.VerifyHeader, err = neofscrypto.SignRequestWithBuffer(neofsecdsa.Signer(signer), req, nil) - }) - if err != nil { - return err - } - - return c.ForAnyGRPCConn(ctx, stream.continueWithConn) - }) - return p, nil -} - func (s *Server) Search(_ *protoobject.SearchRequest, _ protoobject.ObjectService_SearchServer) error { return grpcstatus.Error(grpccodes.Unimplemented, "no longer supported, use SearchV2") } diff --git a/pkg/services/object/server_test.go b/pkg/services/object/server_test.go index a311d80e1f..7aed154cbc 100644 --- a/pkg/services/object/server_test.go +++ b/pkg/services/object/server_test.go @@ -90,10 +90,6 @@ func (x noCallObjectService) Delete(context.Context, deletesvc.Prm) error { panic("must not be called") } -func (x noCallObjectService) GetRange(context.Context, getsvc.RangePrm) error { - panic("must not be called") -} - type noCallTestFSChain struct{} func (*noCallTestFSChain) ForEachContainerNodePublicKey(cid.ID, func([]byte) bool) error { @@ -155,9 +151,6 @@ func (noCallTestReqInfoExtractor) HeadRequestToInfo(context.Context, *protoobjec func (noCallTestReqInfoExtractor) GetRequestToInfo(context.Context, *protoobject.GetRequest, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { panic("must not be called") } -func (noCallTestReqInfoExtractor) RangeRequestToInfo(context.Context, *protoobject.GetRangeRequest, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { - panic("must not be called") -} func (noCallTestReqInfoExtractor) SearchV2RequestToInfo(context.Context, *protoobject.SearchV2Request, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { panic("must not be called") } @@ -199,9 +192,6 @@ func (nopReqInfoExtractor) HeadRequestToInfo(context.Context, *protoobject.HeadR func (nopReqInfoExtractor) GetRequestToInfo(context.Context, *protoobject.GetRequest, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { return v2.RequestInfo{}, nil } -func (nopReqInfoExtractor) RangeRequestToInfo(context.Context, *protoobject.GetRangeRequest, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { - return v2.RequestInfo{}, nil -} func (nopReqInfoExtractor) SearchV2RequestToInfo(context.Context, *protoobject.SearchV2Request, cid.ID, common.RequestTokens) (v2.RequestInfo, error) { return v2.RequestInfo{}, nil }