2020-11-17 12:26:03 +00:00
|
|
|
package engine
|
|
|
|
|
|
|
|
import (
|
2021-05-18 08:12:51 +00:00
|
|
|
"errors"
|
2020-11-17 12:26:03 +00:00
|
|
|
"fmt"
|
|
|
|
|
|
|
|
"github.com/google/uuid"
|
|
|
|
"github.com/nspcc-dev/hrw"
|
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/local_object_storage/shard"
|
2021-11-10 07:08:33 +00:00
|
|
|
"github.com/nspcc-dev/neofs-sdk-go/object"
|
2021-10-08 12:25:45 +00:00
|
|
|
"github.com/panjf2000/ants/v2"
|
2020-11-17 12:26:03 +00:00
|
|
|
)
|
|
|
|
|
2020-11-19 13:04:04 +00:00
|
|
|
var errShardNotFound = errors.New("shard not found")
|
|
|
|
|
2020-11-30 14:58:44 +00:00
|
|
|
type hashedShard struct {
|
|
|
|
sh *shard.Shard
|
|
|
|
}
|
|
|
|
|
2020-11-17 12:26:03 +00:00
|
|
|
// AddShard adds a new shard to the storage engine.
|
|
|
|
//
|
|
|
|
// Returns any error encountered that did not allow adding a shard.
|
|
|
|
// Otherwise returns the ID of the added shard.
|
|
|
|
func (e *StorageEngine) AddShard(opts ...shard.Option) (*shard.ID, error) {
|
|
|
|
e.mtx.Lock()
|
|
|
|
defer e.mtx.Unlock()
|
|
|
|
|
2021-10-08 12:25:45 +00:00
|
|
|
pool, err := ants.NewPool(int(e.shardPoolSize), ants.WithNonblocking(true))
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2022-03-01 08:59:05 +00:00
|
|
|
id, err := generateShardID()
|
|
|
|
if err != nil {
|
|
|
|
return nil, fmt.Errorf("could not generate shard ID: %w", err)
|
|
|
|
}
|
2021-10-08 12:25:45 +00:00
|
|
|
|
2022-03-01 08:59:05 +00:00
|
|
|
sh := shard.New(append(opts,
|
2021-02-17 12:27:40 +00:00
|
|
|
shard.WithID(id),
|
|
|
|
shard.WithExpiredObjectsCallback(e.processExpiredTombstones),
|
|
|
|
)...)
|
2020-11-17 12:26:03 +00:00
|
|
|
|
2022-03-01 08:59:05 +00:00
|
|
|
if err := sh.UpdateID(); err != nil {
|
|
|
|
return nil, fmt.Errorf("could not open shard: %w", err)
|
|
|
|
}
|
|
|
|
|
|
|
|
strID := sh.ID().String()
|
|
|
|
if _, ok := e.shards[strID]; ok {
|
|
|
|
return nil, fmt.Errorf("shard with id %s was already added", strID)
|
|
|
|
}
|
|
|
|
|
|
|
|
e.shards[strID] = sh
|
2021-10-08 12:25:45 +00:00
|
|
|
e.shardPools[strID] = pool
|
|
|
|
|
2022-03-01 08:59:05 +00:00
|
|
|
return sh.ID(), nil
|
2020-11-17 12:26:03 +00:00
|
|
|
}
|
|
|
|
|
2021-04-14 08:47:42 +00:00
|
|
|
func generateShardID() (*shard.ID, error) {
|
2020-11-17 12:26:03 +00:00
|
|
|
uid, err := uuid.NewRandom()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
bin, err := uid.MarshalBinary()
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
return shard.NewIDFromBytes(bin), nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (e *StorageEngine) shardWeight(sh *shard.Shard) float64 {
|
|
|
|
weightValues := sh.WeightValues()
|
|
|
|
|
|
|
|
return float64(weightValues.FreeSpace)
|
|
|
|
}
|
|
|
|
|
2020-11-30 14:58:44 +00:00
|
|
|
func (e *StorageEngine) sortShardsByWeight(objAddr fmt.Stringer) []hashedShard {
|
2020-11-18 12:06:47 +00:00
|
|
|
e.mtx.RLock()
|
|
|
|
defer e.mtx.RUnlock()
|
|
|
|
|
2020-11-30 14:58:44 +00:00
|
|
|
shards := make([]hashedShard, 0, len(e.shards))
|
2020-11-17 12:26:03 +00:00
|
|
|
weights := make([]float64, 0, len(e.shards))
|
|
|
|
|
|
|
|
for _, sh := range e.shards {
|
2020-11-30 14:58:44 +00:00
|
|
|
shards = append(shards, hashedShard{sh})
|
2020-11-17 12:26:03 +00:00
|
|
|
weights = append(weights, e.shardWeight(sh))
|
|
|
|
}
|
|
|
|
|
|
|
|
hrw.SortSliceByWeightValue(shards, weights, hrw.Hash([]byte(objAddr.String())))
|
|
|
|
|
|
|
|
return shards
|
|
|
|
}
|
|
|
|
|
2020-12-01 10:18:25 +00:00
|
|
|
func (e *StorageEngine) unsortedShards() []hashedShard {
|
|
|
|
e.mtx.RLock()
|
|
|
|
defer e.mtx.RUnlock()
|
|
|
|
|
|
|
|
shards := make([]hashedShard, 0, len(e.shards))
|
|
|
|
|
|
|
|
for _, sh := range e.shards {
|
|
|
|
shards = append(shards, hashedShard{sh})
|
|
|
|
}
|
|
|
|
|
|
|
|
return shards
|
|
|
|
}
|
|
|
|
|
|
|
|
func (e *StorageEngine) iterateOverSortedShards(addr *object.Address, handler func(int, *shard.Shard) (stop bool)) {
|
|
|
|
for i, sh := range e.sortShardsByWeight(addr) {
|
|
|
|
if handler(i, sh.sh) {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (e *StorageEngine) iterateOverUnsortedShards(handler func(*shard.Shard) (stop bool)) {
|
|
|
|
for _, sh := range e.unsortedShards() {
|
2020-11-30 14:58:44 +00:00
|
|
|
if handler(sh.sh) {
|
2020-11-17 12:26:03 +00:00
|
|
|
break
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2020-11-19 13:04:04 +00:00
|
|
|
|
|
|
|
// SetShardMode sets mode of the shard with provided identifier.
|
|
|
|
//
|
|
|
|
// Returns an error if shard mode was not set, or shard was not found in storage engine.
|
|
|
|
func (e *StorageEngine) SetShardMode(id *shard.ID, m shard.Mode) error {
|
|
|
|
e.mtx.RLock()
|
|
|
|
defer e.mtx.RUnlock()
|
|
|
|
|
|
|
|
for shID, sh := range e.shards {
|
|
|
|
if id.String() == shID {
|
|
|
|
return sh.SetMode(m)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return errShardNotFound
|
|
|
|
}
|
2020-11-30 14:58:44 +00:00
|
|
|
|
|
|
|
func (s hashedShard) Hash() uint64 {
|
|
|
|
return hrw.Hash(
|
|
|
|
[]byte(s.sh.ID().String()),
|
|
|
|
)
|
|
|
|
}
|