diff --git a/CHANGELOG.md b/CHANGELOG.md index 8bad0083da..7ded9216d3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ Changelog for NeoFS Node ## [Unreleased] ### Added +- Storage node support GET request with missing part index X-Header (#4033) ### Fixed - Session v2 token was not supported in the new container + eACL RPC (#4056) diff --git a/go.sum b/go.sum index abfe7d691f..98517c8464 100644 --- a/go.sum +++ b/go.sum @@ -191,8 +191,6 @@ github.com/nspcc-dev/neofs-api-go/v2 v2.14.1-0.20240827150555-5ce597aa14ea h1:mK github.com/nspcc-dev/neofs-api-go/v2 v2.14.1-0.20240827150555-5ce597aa14ea/go.mod h1:YzhD4EZmC9Z/PNyd7ysC7WXgIgURc9uCG1UWDeV027Y= github.com/nspcc-dev/neofs-contract v0.26.1 h1:7Ii7Q4L3au408LOsIWKiSgfnT1g8G9jo3W7381d41T8= github.com/nspcc-dev/neofs-contract v0.26.1/go.mod h1:pevVF9OWdEN5bweKxOu6ryZv9muCEtS1ppzYM4RfBIo= -github.com/nspcc-dev/neofs-sdk-go v1.0.0-rc.20 h1:8ukaQ4w1eiJd7UnrYo6+DTXB7lwjJFDBLdx/3NdIzDo= -github.com/nspcc-dev/neofs-sdk-go v1.0.0-rc.20/go.mod h1:KtAAolqjEdTNZ8T57FvpWsE/i9/Bpj6wuRPeUs/U85I= github.com/nspcc-dev/neofs-sdk-go v1.0.0-rc.20.0.20260703200507-4f450009c764 h1:7zHfAeRQsWx68owaPuyBSIAs4GIUBL5q6Q335PPGjvo= github.com/nspcc-dev/neofs-sdk-go v1.0.0-rc.20.0.20260703200507-4f450009c764/go.mod h1:MqOYvyQgnXgEYLVzxQcZ12gp/wrJtDqMzMfjaBsisBM= github.com/nspcc-dev/rfc6979 v0.2.4 h1:NBgsdCjhLpEPJZqmC9rciMZDcSY297po2smeaRjw57k= diff --git a/pkg/local_object_storage/engine/ec.go b/pkg/local_object_storage/engine/ec.go index e9bdf2807b..d44924dab7 100644 --- a/pkg/local_object_storage/engine/ec.go +++ b/pkg/local_object_storage/engine/ec.go @@ -54,26 +54,42 @@ func (e *StorageEngine) ReadECPart(_ context.Context, cnr cid.ID, parent oid.ID, // If object is marked as garbage (e.g. via [StorageEngine.MarkGarbage]), // GetECPart returns [apistatus.ErrObjectNotFound]. // +// If allowReturnAnyPart is true, in cases when specified object part ID cannot be +// found, StorageEngine returns any part it has for this EC rule. Returned part ID +// should be found in returned header's attributes. +// // If object is locked (e.g. via [StorageEngine.Lock] or stored locker object), // GetECPart ignores expiration, tombstone and garbage marks. -func (e *StorageEngine) GetECPart(_ context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo) (object.Object, io.ReadCloser, error) { +func (e *StorageEngine) GetECPart(_ context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo, allowReturnAnyPart bool) (object.Object, io.ReadCloser, error) { if e.metrics != nil { defer elapsed(e.metrics.AddGetECPartDuration)() } - var hdr object.Object - var stream io.ReadCloser - return hdr, stream, e.getECPartFunc(cnr, parent, pi, func(s shardInterface, cnr cid.ID, parent oid.ID, pi iec.PartInfo) error { - var err error - hdr, stream, err = s.GetECPart(cnr, parent, pi) - return err - }, func(s shardInterface, partAddr oid.Address) error { - h, str, err := s.GetStream(partAddr, true) - if err == nil { - hdr, stream = *h, str + var ( + hdr object.Object + stream io.ReadCloser + gFunc = func() error { + return e.getECPartFunc(cnr, parent, pi, func(s shardInterface, cnr cid.ID, parent oid.ID, pi iec.PartInfo) error { + var err error + hdr, stream, err = s.GetECPart(cnr, parent, pi) + return err + }, func(s shardInterface, partAddr oid.Address) error { + h, str, err := s.GetStream(partAddr, true) + if err == nil { + hdr, stream = *h, str + } + return err + }) } - return err - }) + ) + + err := gFunc() + if errors.Is(err, apistatus.ErrObjectNotFound) && allowReturnAnyPart { + pi.Index = -1 + return hdr, stream, gFunc() + } + + return hdr, stream, err } func (e *StorageEngine) getECPartFunc(cnr cid.ID, parent oid.ID, pi iec.PartInfo, resolveFn func(shardInterface, cid.ID, oid.ID, iec.PartInfo) error, diff --git a/pkg/local_object_storage/engine/ec_test.go b/pkg/local_object_storage/engine/ec_test.go index 4bd4254763..e76e99f54e 100644 --- a/pkg/local_object_storage/engine/ec_test.go +++ b/pkg/local_object_storage/engine/ec_test.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "io" + "math" "slices" "strconv" "testing" @@ -52,7 +53,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { e := errors.New("any error") require.NoError(t, s.BlockExecution(e)) - _, _, err := s.GetECPart(context.Background(), cnr, parentID, pi) + _, _, err := s.GetECPart(context.Background(), cnr, parentID, pi, false) require.Equal(t, e, err) _, _, err = s.ReadECPart(context.Background(), cnr, parentID, pi, make([]byte, 40<<10)) @@ -84,7 +85,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { s := newEngineWithFixedShardOrder([]shardInterface{shardOK, unimplementedShard{}}) // to ensure 2nd shard is not accessed s.metrics = &m - _, _, _ = s.GetECPart(context.Background(), cnr, parentID, pi) + _, _, _ = s.GetECPart(context.Background(), cnr, parentID, pi, false) require.GreaterOrEqual(t, time.Duration(m.getECPart.Load()), sleepTime) _, _, _ = s.ReadECPart(context.Background(), cnr, parentID, pi, make([]byte, 40<<10)) @@ -102,7 +103,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { s.log = l require.PanicsWithValue(t, "zero object ID returned as error", func() { - _, _, _ = s.GetECPart(context.Background(), cnr, parentID, pi) + _, _, _ = s.GetECPart(context.Background(), cnr, parentID, pi, false) }) lb.AssertEmpty() @@ -140,7 +141,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { } checkOK := func(t *testing.T, s *StorageEngine) { - hdr, rdr, err := s.GetECPart(context.Background(), cnr, parentID, pi) + hdr, rdr, err := s.GetECPart(context.Background(), cnr, parentID, pi, false) require.NoError(t, err) assertGetECPartOK(t, partObj, hdr, rdr) } @@ -155,7 +156,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { require.NoError(t, rdr.Close()) } checkErrorIs := func(t *testing.T, s *StorageEngine, e error) { - _, _, err := s.GetECPart(context.Background(), cnr, parentID, pi) + _, _, err := s.GetECPart(context.Background(), cnr, parentID, pi, false) require.ErrorIs(t, err, e) } checkErrorIsBuffered := func(t *testing.T, s *StorageEngine, e error) { @@ -475,7 +476,7 @@ func TestStorageEngine_GetECPart(t *testing.T) { require.NoError(t, s.Put(context.Background(), &sysObj, nil)) - hdr, rdr, err := s.GetECPart(context.Background(), cnr, sysObj.GetID(), pi) + hdr, rdr, err := s.GetECPart(context.Background(), cnr, sysObj.GetID(), pi, false) require.NoError(t, err) assertGetECPartOK(t, sysObj, hdr, rdr) }) @@ -494,15 +495,51 @@ func TestStorageEngine_GetECPart(t *testing.T) { require.NoError(t, s.Put(context.Background(), &linker, nil)) - hdr, rdr, err := s.GetECPart(context.Background(), cnr, parentID, pi) + hdr, rdr, err := s.GetECPart(context.Background(), cnr, parentID, pi, false) require.NoError(t, err) assertGetECPartOK(t, linker, hdr, rdr) - hdr, rdr, err = s.GetECPart(context.Background(), cnr, linker.GetID(), pi) + hdr, rdr, err = s.GetECPart(context.Background(), cnr, linker.GetID(), pi, false) require.NoError(t, err) assertGetECPartOK(t, linker, hdr, rdr) }) + t.Run("allow any part found", func(t *testing.T) { + piMissing := pi + piMissing.Index++ + piAny := pi + piAny.Index = -1 + piAnother := pi + piAnother.Index = math.MaxInt + + anotherPartObj, err := iec.FormObjectForECPart(neofscryptotest.Signer(), parentObj, testutil.RandByteSlice(32), piAnother) + require.NoError(t, err) + + shardWithObject := &mockShard{ + getECPart: map[getECPartKey]getECPartValue{ + {cnr: cnr, parent: parentID, pi: pi}: {obj: partObj}, + {cnr: cnr, parent: parentID, pi: piMissing}: {err: apistatus.ErrObjectNotFound}, + {cnr: cnr, parent: parentID, pi: piAny}: {obj: anotherPartObj}, + }, + } + shardNoObject := &mockShard{ + getECPart: map[getECPartKey]getECPartValue{ + {cnr: cnr, parent: parentID, pi: pi}: {err: apistatus.ErrObjectNotFound}, + {cnr: cnr, parent: parentID, pi: piMissing}: {err: apistatus.ErrObjectNotFound}, + {cnr: cnr, parent: parentID, pi: piAny}: {err: apistatus.ErrObjectNotFound}, + }, + } + + e := newEngineWithFixedShardOrder([]shardInterface{shardNoObject, shardWithObject, shardNoObject}) + checkOK(t, e) + + _, _, err = e.GetECPart(t.Context(), cnr, parentID, piMissing, false) + require.ErrorIs(t, err, apistatus.ErrObjectNotFound) + + _, _, err = e.GetECPart(t.Context(), cnr, parentID, piMissing, true) + require.NoError(t, err) + }) + l, lb := testutil.NewBufferedLogger(t, zap.DebugLevel) s := newEngineWithFixedShardOrder([]shardInterface{shardOK, unimplementedShard{}}) // to ensure 2nd shard is not accessed @@ -1501,7 +1538,7 @@ func testPutTombstoneEC(t *testing.T) { require.ErrorIs(t, err, target) _, err = s.Head(context.Background(), partAddrs[i], true) require.ErrorIs(t, err, target) - _, _, err = s.GetECPart(context.Background(), cnr, parent.GetID(), iec.PartInfo{RuleIndex: ruleIdx, Index: i}) + _, _, err = s.GetECPart(context.Background(), cnr, parent.GetID(), iec.PartInfo{RuleIndex: ruleIdx, Index: i}, false) require.ErrorIs(t, err, target) } } @@ -1560,7 +1597,7 @@ func testPutTombstoneEC(t *testing.T) { require.NoError(t, err) require.Equal(t, parts[i].CutPayload(), hdr) - h, rdr, err := s.GetECPart(context.Background(), cnr, parent.GetID(), iec.PartInfo{RuleIndex: ruleIdx, Index: i}) + h, rdr, err := s.GetECPart(context.Background(), cnr, parent.GetID(), iec.PartInfo{RuleIndex: ruleIdx, Index: i}, false) assertGetStreamOK(t, &h, rdr, err, parts[i]) } diff --git a/pkg/local_object_storage/metabase/ec.go b/pkg/local_object_storage/metabase/ec.go index aa626d52cf..58d8c380d3 100644 --- a/pkg/local_object_storage/metabase/ec.go +++ b/pkg/local_object_storage/metabase/ec.go @@ -19,6 +19,8 @@ import ( // ResolveECPart resolves object that carries EC part produced within cnr for // parent object and indexed by pi, checks its availability and returns its ID. +// If [iec.PartInfo] has negative part index, resolves the lowest part index it +// stores (if any). // // If the object is not EC part but of [object.TypeTombstone], [object.TypeLock] // or [object.TypeLink] type, ResolveECPart returns its ID instead. @@ -134,6 +136,11 @@ func (db *DB) resolveECPartInMetaBucket(crs *bbolt.Cursor, parent oid.ID, pi iec var rulePref, partPref, typePref []byte var sizeSplitInfo *object.SplitInfo isParent := false + var ( + // for pi.Index < 0 only + foundPartMinInd = -1 + foundPart oid.ID + ) for id := range iterAttrVal(crs, object.FilterParentID, parent[:]) { if id.IsZero() { return oid.ID{}, fmt.Errorf("invalid child of %s parent: %w", parent, oid.ErrZero) @@ -180,14 +187,39 @@ func (db *DB) resolveECPartInMetaBucket(crs *bbolt.Cursor, parent oid.ID, pi iec } if partPref == nil { - partPref = slices.Concat([]byte{metaPrefixIDAttr}, id[:], []byte(iec.AttributePartIdx), objectcore.MetaAttributeDelimiter, []byte(strconv.Itoa(pi.Index))) + partPref = slices.Concat([]byte{metaPrefixIDAttr}, id[:], []byte(iec.AttributePartIdx), objectcore.MetaAttributeDelimiter) + if pi.Index >= 0 { + partPref = append(partPref, strconv.Itoa(pi.Index)...) + } } else { copy(partPref[1:], id[:]) } - if k, _ = partCrs.Seek(partPref); bytes.Equal(k, partPref) { + k, _ = partCrs.Seek(partPref) + if !bytes.HasPrefix(k, partPref) { + continue + } + if pi.Index >= 0 { return id, nil } + + ind, err := strconv.Atoi(string(k[len(partPref):])) + if err != nil { + panic(fmt.Errorf("unexpected EC part index from %X value: %w", k, err)) + } + if foundPartMinInd == -1 { + foundPartMinInd = ind + foundPart = id + continue + } + if ind < foundPartMinInd { + foundPartMinInd = ind + foundPart = id + } } + if pi.Index < 0 && !foundPart.IsZero() { + return foundPart, nil + } + if !isParent { // neither tombstone nor lock can be a parent if typePref == nil { typePref = make([]byte, metaIDTypePrefixSize) diff --git a/pkg/local_object_storage/metabase/ec_test.go b/pkg/local_object_storage/metabase/ec_test.go index 666d9e17f0..9ab5c36cfc 100644 --- a/pkg/local_object_storage/metabase/ec_test.go +++ b/pkg/local_object_storage/metabase/ec_test.go @@ -497,6 +497,28 @@ func TestDB_ResolveECPart(t *testing.T) { }) }) + t.Run("negative part index", func(t *testing.T) { + db := newDB(t) + + const ruleInd = 1234 + pi1 := iec.PartInfo{ + RuleIndex: ruleInd, + Index: 0, + } + pi2 := iec.PartInfo{ + RuleIndex: ruleInd, + Index: 1, + } + + p1 := addPart(t, db, pi1) + p2 := addPart(t, db, pi2) + + checkOK(t, db, pi1, p1) + checkOK(t, db, pi2, p2) + + checkOKWithLen(t, db, iec.PartInfo{RuleIndex: ruleInd, Index: -1}, p1, 1000+pi1.Index) + }) + db := newDB(t) require.NoError(t, db.Put(&partObj)) checkOK(t, db, pi, partID) diff --git a/pkg/local_object_storage/shard/ec.go b/pkg/local_object_storage/shard/ec.go index 477594217b..d9e149bee4 100644 --- a/pkg/local_object_storage/shard/ec.go +++ b/pkg/local_object_storage/shard/ec.go @@ -36,6 +36,8 @@ func (s *Shard) ReadECPart(cnr cid.ID, parent oid.ID, pi iec.PartInfo, buf []byt // parent object and indexed by pi in the underlying metabase, checks its // availability and reads it from the underlying BLOB storage. The result is a // header and a payload stream that must be closed by caller after processing. +// If [iec.PartInfo] has negative part index, returned the lowest part index it +// stores (if any). // // If the object is not EC part but of [object.TypeTombstone], [object.TypeLock] // or [object.TypeLink] type, GetECPart this object instead. diff --git a/pkg/services/object/get/ec.go b/pkg/services/object/get/ec.go index 0b963e3270..d18f6ac8e8 100644 --- a/pkg/services/object/get/ec.go +++ b/pkg/services/object/get/ec.go @@ -21,6 +21,7 @@ import ( "github.com/nspcc-dev/neofs-node/pkg/services/object/internal" apistatus "github.com/nspcc-dev/neofs-sdk-go/client/status" "github.com/nspcc-dev/neofs-sdk-go/container" + "github.com/nspcc-dev/neofs-sdk-go/container/acl" 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" @@ -41,8 +42,8 @@ func (x tooManyPartsUnavailableError) Error() string { // // Returns [apistatus.ErrObjectAlreadyRemoved] if the object was marked for // removal. Returns [apistatus.ErrObjectNotFound] if the object is missing. -func (s *Service) copyLocalECPart(ctx context.Context, dst ObjectWriter, cnr cid.ID, parent oid.ID, pi iec.PartInfo) error { - hdr, rc, err := s.localObjects.GetECPart(ctx, cnr, parent, pi) +func (s *Service) copyLocalECPart(ctx context.Context, dst ObjectWriter, cnr cid.ID, parent oid.ID, pi iec.PartInfo, allowAnyPart bool) error { + hdr, rc, err := s.localObjects.GetECPart(ctx, cnr, parent, pi, allowAnyPart) if err != nil { return fmt.Errorf("get object from local storage: %w", err) } @@ -528,7 +529,7 @@ func (s *Service) getECPartStream(ctx context.Context, cnr cid.ID, parent oid.ID } if local { - partHdr, rc, err = s.localObjects.GetECPart(ctx, cnr, parent, pi) + partHdr, rc, err = s.localObjects.GetECPart(ctx, cnr, parent, pi, false) } else { partHdr, rc, err = s.getECPartFromNode(ctx, cnr, parent, sTok, pi, sortedNodes[i]) } @@ -1636,10 +1637,7 @@ func (s *Service) streamECPartRangePrefix(ctx context.Context, transport GetECRe return copiedLen, nil } -// returns [iec.PartInfo.RuleIndex] = -1 if request is not for particular EC part. -func checkECPartInfoRequest(xHdrs []string, cnr container.Container) (iec.PartInfo, error) { - var res iec.PartInfo - +func findECIndexes(xHdrs []string) (string, string) { var ruleIdxStr, partIdxStr string for i := 0; i < len(xHdrs); i += 2 { if xHdrs[i] == iec.AttributeRuleIdx { @@ -1653,10 +1651,22 @@ func checkECPartInfoRequest(xHdrs []string, cnr container.Container) (iec.PartIn } } + return ruleIdxStr, partIdxStr +} + +// returns [iec.PartInfo.RuleIndex] = -1 if request is not for particular EC part. +func checkECPartInfoRequest(xHdrs []string, cnr container.Container) (iec.PartInfo, error) { + var ( + res iec.PartInfo + ruleIdxStr, partIdxStr = findECIndexes(xHdrs) + bAcl = cnr.BasicACL() + ) + if ruleIdxStr == "" && partIdxStr == "" { res.RuleIndex = -1 return res, nil } + var ecRules = cnr.PlacementPolicy().ECRules() if (ruleIdxStr == "") != (partIdxStr == "") { return res, fmt.Errorf("%s and %s X-headers must be set together", iec.AttributeRuleIdx, iec.AttributePartIdx) @@ -1667,34 +1677,117 @@ func checkECPartInfoRequest(xHdrs []string, cnr container.Container) (iec.PartIn if err != nil { return res, fmt.Errorf("invalid %s X-header: %w", iec.AttributeRuleIdx, err) } + res.RuleIndex = int(ruleIdx) + + err = checkECRuleIdx(bAcl, ecRules, res.RuleIndex) + if err != nil { + return res, err + } partIdx, err := strconv.ParseUint(partIdxStr, 10, 8) if err != nil { return res, fmt.Errorf("invalid %s X-header: %w", iec.AttributePartIdx, err) } + res.Index = int(partIdx) - if cnr.BasicACL() != 0 { // Uninitialized in tests, safe to do anyway, invalid requests will fail. - var ecRules = cnr.PlacementPolicy().ECRules() + return res, checkECPartIdx(bAcl, cnr.PlacementPolicy().ECRules(), res.RuleIndex, res.Index) +} - if len(ecRules) == 0 { - return res, errors.New("EC part requested in container without EC policy") - } +// returns [iec.PartInfo.RuleIndex] = -1 if request is not for particular EC part. +func checkECPartInfoGetRequest(neofs NeoFSNetwork, prm Prm) (iec.PartInfo, error) { + var ( + res iec.PartInfo + xHdrs = prm.common.XHeaders() + cnr = prm.container + ruleIdxStr, partIdxStr = findECIndexes(xHdrs) + ) - if int(ruleIdx) >= len(ecRules) { - return res, fmt.Errorf("EC rule index overflows container policy: idx=%d,rules=%d", ruleIdx, len(ecRules)) + if ruleIdxStr == "" { + if partIdxStr != "" { + return res, fmt.Errorf("%s X-header must be set in a correct EC part GET request (%s X-header found: %s)", + iec.AttributeRuleIdx, iec.AttributePartIdx, partIdxStr) } - if total := ecRules[ruleIdx].DataPartNum() + ecRules[ruleIdx].ParityPartNum(); int(partIdx) >= int(total) { - return res, fmt.Errorf("EC part index overflows container policy: idx=%d,parts=%d", partIdx, total) - } + res.RuleIndex = -1 + return res, nil + } else if partIdxStr == "" && neofs == nil { + return res, fmt.Errorf("request must have %s header for EC objects", iec.AttributePartIdx) } + var ( + bACL = cnr.BasicACL() + ecRules = cnr.PlacementPolicy().ECRules() + ) + + // TODO: state limits in https://github.com/nspcc-dev/neofs-api. Share consts for them. + ruleIdx, err := strconv.ParseUint(ruleIdxStr, 10, 8) + if err != nil { + return res, fmt.Errorf("invalid %s X-header: %w", iec.AttributeRuleIdx, err) + } res.RuleIndex = int(ruleIdx) + err = checkECRuleIdx(bACL, ecRules, res.RuleIndex) + if err != nil { + return res, err + } + + if partIdxStr == "" { + nodes, repList, ecList, err := neofs.GetNodesForObject(prm.addr) + if err != nil { + return res, fmt.Errorf("missing %s and calculating node's EC part finished with error: %w", iec.AttributePartIdx, err) + } + if res.RuleIndex >= len(ecList) { + return res, fmt.Errorf("%s X-header is bigger than placement rules number (%d >= %d)", iec.AttributeRuleIdx, res.RuleIndex, len(ecList)) + } + nodesForObject := nodes[len(repList)+res.RuleIndex] + i := slices.IndexFunc(nodesForObject, func(info netmap.NodeInfo) bool { + return neofs.IsLocalNodePublicKey(info.PublicKey()) + }) + if i == -1 { + return res, fmt.Errorf("missing %s and node does not belong to placement list for %d EC rule", iec.AttributePrefix, res.RuleIndex) + } + + res.Index = i + + return res, nil + } + + partIdx, err := strconv.ParseUint(partIdxStr, 10, 8) + if err != nil { + return res, fmt.Errorf("invalid %s X-header: %w", iec.AttributePartIdx, err) + } res.Index = int(partIdx) + err = checkECPartIdx(bACL, ecRules, res.RuleIndex, res.Index) + if err != nil { + return res, err + } + return res, nil } +func checkECRuleIdx(bACL acl.Basic, ecRules []netmap.ECRule, ruleIdx int) error { + if bACL != 0 { // Uninitialized in tests, safe to do anyway, invalid requests will fail. + if len(ecRules) == 0 { + return errors.New("EC part requested in container without EC policy") + } + if ruleIdx >= len(ecRules) { + return fmt.Errorf("EC rule index overflows container policy: idx=%d,rules=%d", ruleIdx, len(ecRules)) + } + } + + return nil +} + +func checkECPartIdx(bACL acl.Basic, ecRules []netmap.ECRule, ruleIdx, partIdx int) error { + if bACL != 0 { // Uninitialized in tests, safe to do anyway, invalid requests will fail. + if total := ecRules[ruleIdx].DataPartNum() + ecRules[ruleIdx].ParityPartNum(); partIdx >= int(total) { + return fmt.Errorf("EC part index overflows container policy: idx=%d,parts=%d", partIdx, total) + } + } + + return nil +} + func checkECAttributesInReceivedObject(hdr object.Object, ruleIdx, partIdx string) error { var found uint8 const expected = 2 diff --git a/pkg/services/object/get/ec_test.go b/pkg/services/object/get/ec_test.go index 4309bae3d7..f4c901d7e0 100644 --- a/pkg/services/object/get/ec_test.go +++ b/pkg/services/object/get/ec_test.go @@ -42,20 +42,15 @@ func TestService_Get_EC_Part(t *testing.T) { xs []string assertErr func(t *testing.T, err error) }{ - {name: "rule idx only", xs: []string{ - "__NEOFS__EC_RULE_IDX", "0", - }, assertErr: func(t *testing.T, err error) { - require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX and __NEOFS__EC_PART_IDX X-headers must be set together") - }}, {name: "part idx only", xs: []string{ "__NEOFS__EC_PART_IDX", "0", }, assertErr: func(t *testing.T, err error) { - require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX and __NEOFS__EC_PART_IDX X-headers must be set together") + require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX X-header must be set in a correct EC part GET request") }}, {name: "empty rule idx", xs: []string{ "__NEOFS__EC_RULE_IDX", "", "__NEOFS__EC_PART_IDX", "0", }, assertErr: func(t *testing.T, err error) { - require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX and __NEOFS__EC_PART_IDX X-headers must be set together") + require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX X-header must be set in a correct EC part GET request") }}, {name: "non-numeric rule idx", xs: []string{ "__NEOFS__EC_RULE_IDX", "foo", "__NEOFS__EC_PART_IDX", "0", @@ -81,11 +76,6 @@ func TestService_Get_EC_Part(t *testing.T) { require.ErrorContains(t, err, "invalid __NEOFS__EC_RULE_IDX X-header") require.ErrorContains(t, err, "value out of range") }}, - {name: "empty part idx", xs: []string{ - "__NEOFS__EC_RULE_IDX", "0", "__NEOFS__EC_PART_IDX", "", - }, assertErr: func(t *testing.T, err error) { - require.ErrorContains(t, err, "__NEOFS__EC_RULE_IDX and __NEOFS__EC_PART_IDX X-headers must be set together") - }}, {name: "non-numeric part idx", xs: []string{ "__NEOFS__EC_RULE_IDX", "0", "__NEOFS__EC_PART_IDX", "foo", }, assertErr: func(t *testing.T, err error) { @@ -306,6 +296,19 @@ func TestService_Get_EC_Part(t *testing.T) { require.EqualValues(t, 1, closer.count.Load()) }) + t.Run("placement build failure", func(t *testing.T) { + placementErr := errors.New("some placement error") + svc := New(&mockNeoFSNet{err: placementErr}) + + var prm Prm + prm.WithAddress(parentAddr) + parameterizePartInfoString(t, &prm, "1", "") + + err := svc.Get(ctx, prm) + require.ErrorIs(t, err, placementErr) + require.EqualError(t, err, fmt.Sprintf("invalid request: missing %s and calculating node's EC part finished with error: %s", iec.AttributePartIdx, placementErr)) + }) + svc := New(&mockNeoFSNet{ getNodesForObject: map[oid.Address]getNodesForObjectValue{ parentAddr: {ecRules: ecRules}, diff --git a/pkg/services/object/get/get.go b/pkg/services/object/get/get.go index 2edd80a499..b5fafd90a0 100644 --- a/pkg/services/object/get/get.go +++ b/pkg/services/object/get/get.go @@ -15,7 +15,13 @@ import ( // Get serves a request to get an object by address, and returns Streamer instance. func (s *Service) Get(ctx context.Context, prm Prm) error { - pi, err := checkECPartInfoRequest(prm.common.XHeaders(), prm.container) + var neofsNet NeoFSNetwork + // range requests do not support fetching additional info about EC parts + if prm.rng == nil { + neofsNet = s.neoFSNet + } + + pi, err := checkECPartInfoGetRequest(neofsNet, prm) if err != nil { // TODO: track https://github.com/nspcc-dev/neofs-api/issues/269. return fmt.Errorf("invalid request: %w", err) @@ -43,7 +49,7 @@ func (s *Service) Get(ctx context.Context, prm Prm) error { return err } - return s.copyLocalECPart(ctx, prm.objWriter, prm.addr.Container(), prm.addr.Object(), pi) + return s.copyLocalECPart(ctx, prm.objWriter, prm.addr.Container(), prm.addr.Object(), pi, prm.ecReturnAnyPart) } if prm.common.LocalOnly() && diff --git a/pkg/services/object/get/prm.go b/pkg/services/object/get/prm.go index e7fbf94311..fd7b2b695c 100644 --- a/pkg/services/object/get/prm.go +++ b/pkg/services/object/get/prm.go @@ -32,7 +32,8 @@ type Prm struct { localGetBuffer []byte submitLocalGetStreamFn SubmitStreamFunc - ecTransport GetECRequestTransport + ecTransport GetECRequestTransport + ecReturnAnyPart bool transportFn GetTransportFunc } @@ -277,6 +278,12 @@ func (p *Prm) WithECTransport(transport GetECRequestTransport) { p.ecTransport = transport } +// WithECReturnAnyPart makes return any EC part local node has, +// if recuested part ID is not found. +func (p *Prm) WithECReturnAnyPart() { + p.ecReturnAnyPart = true +} + // SetForwardRequestFunc specifies transport implementation for request // forwarding to container nodes by non-container server. func (p *commonPrm) SetForwardRequestFunc(f ForwardRequestFunc) { diff --git a/pkg/services/object/get/service.go b/pkg/services/object/get/service.go index 4da9853f5e..d1f3614309 100644 --- a/pkg/services/object/get/service.go +++ b/pkg/services/object/get/service.go @@ -142,7 +142,7 @@ type cfg struct { // // Returns [apistatus.ErrObjectAlreadyRemoved] if the object was marked for // removal. Returns [apistatus.ErrObjectNotFound] if the object is missing. - GetECPart(ctx context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo) (object.Object, io.ReadCloser, error) + GetECPart(ctx context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo, allowAnyPart bool) (object.Object, io.ReadCloser, error) // GetECPartRange reads specified payload ranage of stored object that carries // EC part produced within cnr for parent object and indexed by pi. // diff --git a/pkg/services/object/get/service_test.go b/pkg/services/object/get/service_test.go index a49ab54e6e..879703f3d9 100644 --- a/pkg/services/object/get/service_test.go +++ b/pkg/services/object/get/service_test.go @@ -94,11 +94,16 @@ type getNodesForObjectValue struct { } type mockNeoFSNet struct { + err error getNodesForObject map[oid.Address]getNodesForObjectValue localPubKey []byte } func (x *mockNeoFSNet) GetNodesForObject(addr oid.Address) ([][]netmap.NodeInfo, []uint, []iec.Rule, error) { + if x.err != nil { + return nil, nil, nil, x.err + } + v, ok := x.getNodesForObject[addr] if !ok { return nil, nil, nil, errors.New("[test] unexpected object requested") @@ -127,7 +132,7 @@ type mockLocalObjects struct { getECPart map[getECPartKey]getECPartValue } -func (x *mockLocalObjects) GetECPart(_ context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo) (object.Object, io.ReadCloser, error) { +func (x *mockLocalObjects) GetECPart(_ context.Context, cnr cid.ID, parent oid.ID, pi iec.PartInfo, allowAnyPart bool) (object.Object, io.ReadCloser, error) { v, ok := x.getECPart[getECPartKey{ cnr: cnr, parent: parent, @@ -188,7 +193,7 @@ func (x unimplementedLocalStorage) GetECPartRange(_ context.Context, _ cid.ID, _ panic("unimplemented") } -func (unimplementedLocalStorage) GetECPart(_ context.Context, _ cid.ID, _ oid.ID, _ iec.PartInfo) (object.Object, io.ReadCloser, error) { +func (unimplementedLocalStorage) GetECPart(_ context.Context, _ cid.ID, _ oid.ID, _ iec.PartInfo, _ bool) (object.Object, io.ReadCloser, error) { panic("unimplemented") } diff --git a/pkg/services/object/server.go b/pkg/services/object/server.go index b22e454dc9..9941472bbc 100644 --- a/pkg/services/object/server.go +++ b/pkg/services/object/server.go @@ -920,6 +920,9 @@ type getStream struct { recheckEACL bool signResponse bool payloadOnly bool + + sendECPartIndInResponse bool + ecFoundPartInd string } func (s *getStream) ValidateHeader(hdr *object.Object) error { @@ -949,6 +952,14 @@ func (s *getStream) WriteHeader(hdr *object.Object) error { if err := s.ValidateHeader(hdr); err != nil { return err } + if s.sendECPartIndInResponse { + for _, attr := range hdr.Attributes() { + if attr.Key() == iec.AttributePartIdx { + s.ecFoundPartInd = attr.Value() + break + } + } + } if s.payloadOnly { return nil } @@ -967,13 +978,24 @@ func (s *getStream) WriteHeader(hdr *object.Object) error { } func (s *getStream) WriteChunk(chunk []byte) error { - for buf := bytes.NewBuffer(chunk); buf.Len() > 0; { + var metaHeader *protosession.ResponseMetaHeader + if s.ecFoundPartInd != "" { + metaHeader = &protosession.ResponseMetaHeader{ + XHeaders: []*protosession.XHeader{{ + Key: iec.AttributePartIdx, + Value: s.ecFoundPartInd, + }}, + } + } + + for buf := bytes.NewBuffer(chunk); buf.Len() > 0; metaHeader, s.ecFoundPartInd = nil, "" { newResp := &protoobject.GetResponse{ Body: &protoobject.GetResponse_Body{ ObjectPart: &protoobject.GetResponse_Body_Chunk{ Chunk: buf.Next(maxRespDataChunkSize), }, }, + MetaHeader: metaHeader, } if err := s.srv.sendGetResponse(s.base, newResp, s.signResponse); err != nil { return err @@ -1039,14 +1061,15 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ } p, err := convertGetPrm(s.signer, reqInfo.Container, req, &getStream{ - base: gStream, - srv: s, - reqCID: cnrID, - reqOID: objID, - reqInfo: reqInfo, - recheckEACL: recheckEACL, - signResponse: needSignResp, - payloadOnly: req.GetBody().GetPayloadOnly(), + base: gStream, + srv: s, + reqCID: cnrID, + reqOID: objID, + reqInfo: reqInfo, + recheckEACL: recheckEACL, + signResponse: needSignResp, + payloadOnly: req.GetBody().GetPayloadOnly(), + sendECPartIndInResponse: sendECPartIdxInResponse(req), }, cnrID, objID, reqMD) if err != nil { if !errors.Is(err, apistatus.Error) { @@ -1127,6 +1150,31 @@ func (s *Server) Get(req *protoobject.GetRequest, gStream protoobject.ObjectServ return nil } +func sendECPartIdxInResponse(req *protoobject.GetRequest) bool { + if req.MetaHeader == nil || !req.Body.PayloadOnly || req.Body.Range != nil { + return false + } + var ( + isECPartReq bool + partIdxFound bool + ) +attrL: + for _, attr := range req.MetaHeader.XHeaders { + switch attr.Key { + case iec.AttributePartIdx: + isECPartReq = true + partIdxFound = attr.Key != "" + break attrL + case iec.AttributeRuleIdx: + isECPartReq = true + default: + continue + } + } + + return isECPartReq && !partIdxFound +} + func (s *Server) copyGetStream(gStream grpc.ServerStream, hdrRespBuf *iprotobuf.MemBuffer, hdrBuf []byte, prefixLen, hdrTo int, stream io.Reader, pldFldOff int, needSignResp bool) error { var (