Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
2 changes: 0 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
42 changes: 29 additions & 13 deletions pkg/local_object_storage/engine/ec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Comment thread
roman-khimov marked this conversation as resolved.
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
Comment thread
roman-khimov marked this conversation as resolved.
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,
Expand Down
57 changes: 47 additions & 10 deletions pkg/local_object_storage/engine/ec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"errors"
"fmt"
"io"
"math"
"slices"
"strconv"
"testing"
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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))
Expand All @@ -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()
Expand Down Expand Up @@ -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)
}
Expand All @@ -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) {
Expand Down Expand Up @@ -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)
})
Expand All @@ -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
Expand Down Expand Up @@ -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)
}
}
Expand Down Expand Up @@ -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])
}

Expand Down
36 changes: 34 additions & 2 deletions pkg/local_object_storage/metabase/ec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
22 changes: 22 additions & 0 deletions pkg/local_object_storage/metabase/ec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
2 changes: 2 additions & 0 deletions pkg/local_object_storage/shard/ec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading