From bb74e97958244a3286f9e6743c4a2fa9b518f07e Mon Sep 17 00:00:00 2001 From: Pavel Karpy Date: Thu, 2 Jul 2026 20:24:29 +0300 Subject: [PATCH 1/2] sn: allow EC GET with missing part ID If X-Headers do not include EC part ID, SN fills it with an ID that it should store according to current Netmap. It such part cannot be found, it will return first part ID it has locally. This does not relate any form of GETRAGE reques, both GET with RANGE and native GETRAGE. This commit implements https://github.com/nspcc-dev/neofs-api/pull/393. Signed-off-by: Pavel Karpy --- CHANGELOG.md | 1 + pkg/local_object_storage/engine/ec.go | 42 ++++-- pkg/local_object_storage/engine/ec_test.go | 57 +++++++-- pkg/local_object_storage/metabase/ec.go | 36 +++++- pkg/local_object_storage/metabase/ec_test.go | 22 ++++ pkg/local_object_storage/shard/ec.go | 2 + pkg/services/object/get/ec.go | 127 ++++++++++++++++--- pkg/services/object/get/ec_test.go | 27 ++-- pkg/services/object/get/get.go | 10 +- pkg/services/object/get/prm.go | 9 +- pkg/services/object/get/service.go | 2 +- pkg/services/object/get/service_test.go | 9 +- pkg/services/object/server.go | 66 ++++++++-- 13 files changed, 341 insertions(+), 69 deletions(-) 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/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 ( From c7dac28309202f4c1f412677804086c1fb7651c8 Mon Sep 17 00:00:00 2001 From: Pavel Karpy Date: Tue, 7 Jul 2026 16:07:33 +0300 Subject: [PATCH 2/2] go.sum: cleanup old SDK hashes Signed-off-by: Pavel Karpy --- go.sum | 2 -- 1 file changed, 2 deletions(-) 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=