2020-09-28 13:22:13 +00:00
|
|
|
package rangehashsvc
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"crypto/sha256"
|
|
|
|
"fmt"
|
|
|
|
|
|
|
|
"github.com/nspcc-dev/neofs-api-go/pkg"
|
2020-11-23 12:59:06 +00:00
|
|
|
"github.com/nspcc-dev/neofs-api-go/pkg/client"
|
2020-09-28 13:22:13 +00:00
|
|
|
"github.com/nspcc-dev/neofs-api-go/pkg/object"
|
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/core/container"
|
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/core/netmap"
|
2020-11-19 08:27:40 +00:00
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/local_object_storage/engine"
|
2020-09-28 13:22:13 +00:00
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/network"
|
2020-11-18 13:04:59 +00:00
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/network/cache"
|
2020-12-07 17:49:47 +00:00
|
|
|
getsvc "github.com/nspcc-dev/neofs-node/pkg/services/object/get"
|
2020-09-28 13:22:13 +00:00
|
|
|
headsvc "github.com/nspcc-dev/neofs-node/pkg/services/object/head"
|
|
|
|
objutil "github.com/nspcc-dev/neofs-node/pkg/services/object/util"
|
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/util"
|
2020-11-23 11:24:51 +00:00
|
|
|
"github.com/nspcc-dev/neofs-node/pkg/util/logger"
|
2020-09-28 13:22:13 +00:00
|
|
|
"github.com/pkg/errors"
|
2020-11-23 11:24:51 +00:00
|
|
|
"go.uber.org/zap"
|
2020-09-28 13:22:13 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
type Service struct {
|
|
|
|
*cfg
|
|
|
|
}
|
|
|
|
|
|
|
|
type Option func(*cfg)
|
|
|
|
|
|
|
|
type cfg struct {
|
2020-09-29 16:44:59 +00:00
|
|
|
keyStorage *objutil.KeyStorage
|
2020-09-28 13:22:13 +00:00
|
|
|
|
2020-11-19 08:27:40 +00:00
|
|
|
localStore *engine.StorageEngine
|
2020-09-28 13:22:13 +00:00
|
|
|
|
|
|
|
cnrSrc container.Source
|
|
|
|
|
|
|
|
netMapSrc netmap.Source
|
|
|
|
|
|
|
|
workerPool util.WorkerPool
|
|
|
|
|
|
|
|
localAddrSrc network.LocalAddressSource
|
|
|
|
|
|
|
|
headSvc *headsvc.Service
|
|
|
|
|
2020-12-07 17:49:47 +00:00
|
|
|
rangeSvc *getsvc.Service
|
2020-11-18 13:04:59 +00:00
|
|
|
|
|
|
|
clientCache *cache.ClientCache
|
2020-11-23 11:24:51 +00:00
|
|
|
|
|
|
|
log *logger.Logger
|
2020-11-23 12:59:06 +00:00
|
|
|
|
|
|
|
clientOpts []client.Option
|
2020-09-28 13:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func defaultCfg() *cfg {
|
|
|
|
return &cfg{
|
|
|
|
workerPool: new(util.SyncWorkerPool),
|
2020-11-23 11:24:51 +00:00
|
|
|
log: zap.L(),
|
2020-09-28 13:22:13 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewService(opts ...Option) *Service {
|
|
|
|
c := defaultCfg()
|
|
|
|
|
|
|
|
for i := range opts {
|
|
|
|
opts[i](c)
|
|
|
|
}
|
|
|
|
|
|
|
|
return &Service{
|
|
|
|
cfg: c,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *Service) GetRangeHash(ctx context.Context, prm *Prm) (*Response, error) {
|
|
|
|
headResult, err := s.headSvc.Head(ctx, new(headsvc.Prm).
|
|
|
|
WithAddress(prm.addr).
|
2020-09-29 15:05:22 +00:00
|
|
|
WithCommonPrm(prm.common),
|
2020-09-28 13:22:13 +00:00
|
|
|
)
|
|
|
|
if err != nil {
|
|
|
|
return nil, errors.Wrapf(err, "(%T) could not receive Head result", s)
|
|
|
|
}
|
|
|
|
|
|
|
|
origin := headResult.Header()
|
|
|
|
|
2020-11-16 09:43:52 +00:00
|
|
|
originSize := origin.PayloadSize()
|
2020-09-28 13:22:13 +00:00
|
|
|
|
|
|
|
var minLeft, maxRight uint64
|
|
|
|
for i := range prm.rngs {
|
|
|
|
left := prm.rngs[i].GetOffset()
|
|
|
|
right := left + prm.rngs[i].GetLength()
|
|
|
|
|
|
|
|
if originSize < right {
|
|
|
|
return nil, errors.Errorf("(%T) requested payload range is out-of-bounds", s)
|
|
|
|
}
|
|
|
|
|
|
|
|
if left < minLeft {
|
|
|
|
minLeft = left
|
|
|
|
}
|
|
|
|
|
|
|
|
if right > maxRight {
|
|
|
|
maxRight = right
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
borderRng := new(object.Range)
|
|
|
|
borderRng.SetOffset(minLeft)
|
|
|
|
borderRng.SetLength(maxRight - minLeft)
|
|
|
|
|
2020-12-05 12:29:05 +00:00
|
|
|
return s.getHashes(ctx, prm, objutil.NewRangeTraverser(originSize, origin, borderRng))
|
2020-09-28 13:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (s *Service) getHashes(ctx context.Context, prm *Prm, traverser *objutil.RangeTraverser) (*Response, error) {
|
|
|
|
addr := object.NewAddress()
|
2020-11-16 09:43:52 +00:00
|
|
|
addr.SetContainerID(prm.addr.ContainerID())
|
2020-09-28 13:22:13 +00:00
|
|
|
|
|
|
|
resp := &Response{
|
|
|
|
hashes: make([][]byte, 0, len(prm.rngs)),
|
|
|
|
}
|
|
|
|
|
|
|
|
for _, rng := range prm.rngs {
|
|
|
|
for {
|
|
|
|
nextID, nextRng := traverser.Next()
|
|
|
|
if nextRng != nil {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
|
|
|
addr.SetObjectID(nextID)
|
|
|
|
|
|
|
|
head, err := s.headSvc.Head(ctx, new(headsvc.Prm).
|
|
|
|
WithAddress(addr).
|
2020-09-29 15:05:22 +00:00
|
|
|
WithCommonPrm(prm.common),
|
2020-09-28 13:22:13 +00:00
|
|
|
)
|
|
|
|
if err != nil {
|
|
|
|
return nil, errors.Wrapf(err, "(%T) could not receive object header", s)
|
|
|
|
}
|
|
|
|
|
|
|
|
traverser.PushHeader(head.Header())
|
|
|
|
}
|
|
|
|
|
|
|
|
traverser.SetSeekRange(rng)
|
|
|
|
|
|
|
|
var hasher hasher
|
|
|
|
|
|
|
|
for {
|
|
|
|
nextID, nextRng := traverser.Next()
|
|
|
|
|
|
|
|
if hasher == nil {
|
|
|
|
if nextRng.GetLength() == rng.GetLength() {
|
|
|
|
hasher = new(singleHasher)
|
|
|
|
} else {
|
|
|
|
switch prm.typ {
|
|
|
|
default:
|
|
|
|
panic(fmt.Sprintf("unexpected checksum type %v", prm.typ))
|
|
|
|
case pkg.ChecksumSHA256:
|
|
|
|
hasher = &commonHasher{h: sha256.New()}
|
|
|
|
case pkg.ChecksumTZ:
|
|
|
|
hasher = &tzHasher{
|
|
|
|
hashes: make([][]byte, 0, 10),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
if nextRng.GetLength() == 0 {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
|
|
|
addr.SetObjectID(nextID)
|
|
|
|
|
|
|
|
if prm.typ == pkg.ChecksumSHA256 && nextRng.GetLength() != rng.GetLength() {
|
|
|
|
// here we cannot receive SHA256 checksum through GetRangeHash service
|
|
|
|
// since SHA256 is not homomorphic
|
2020-12-07 17:49:47 +00:00
|
|
|
rngPrm := getsvc.RangePrm{}
|
|
|
|
rngPrm.SetRange(nextRng)
|
|
|
|
rngPrm.WithAddress(addr)
|
|
|
|
rngPrm.SetChunkWriter(hasher)
|
|
|
|
|
|
|
|
err := s.rangeSvc.GetRange(ctx, rngPrm)
|
2020-09-28 13:22:13 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, errors.Wrapf(err, "(%T) could not receive payload range for %v checksum", s, prm.typ)
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
resp, err := (&distributedHasher{
|
|
|
|
cfg: s.cfg,
|
|
|
|
}).head(ctx, new(Prm).
|
|
|
|
WithAddress(addr).
|
|
|
|
WithChecksumType(prm.typ).
|
2020-09-29 15:05:22 +00:00
|
|
|
FromRanges(nextRng).
|
|
|
|
WithCommonPrm(prm.common),
|
2020-09-28 13:22:13 +00:00
|
|
|
)
|
|
|
|
if err != nil {
|
|
|
|
return nil, errors.Wrapf(err, "(%T) could not receive %v checksum", s, prm.typ)
|
|
|
|
}
|
|
|
|
|
|
|
|
hs := resp.Hashes()
|
|
|
|
if ln := len(hs); ln != 1 {
|
|
|
|
return nil, errors.Errorf("(%T) unexpected %v hashes amount %d", s, prm.typ, ln)
|
|
|
|
}
|
|
|
|
|
2020-12-07 17:49:47 +00:00
|
|
|
_ = hasher.WriteChunk(hs[0])
|
2020-09-28 13:22:13 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
traverser.PushSuccessSize(nextRng.GetLength())
|
|
|
|
}
|
|
|
|
|
|
|
|
sum, err := hasher.sum()
|
|
|
|
if err != nil {
|
|
|
|
return nil, errors.Wrapf(err, "(%T) could not calculate %v checksum", s, prm.typ)
|
|
|
|
}
|
|
|
|
|
|
|
|
resp.hashes = append(resp.hashes, sum)
|
|
|
|
}
|
|
|
|
|
|
|
|
return resp, nil
|
|
|
|
}
|
|
|
|
|
2020-09-29 16:44:59 +00:00
|
|
|
func WithKeyStorage(v *objutil.KeyStorage) Option {
|
2020-09-28 13:22:13 +00:00
|
|
|
return func(c *cfg) {
|
2020-09-29 16:44:59 +00:00
|
|
|
c.keyStorage = v
|
2020-09-28 13:22:13 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-11-19 08:27:40 +00:00
|
|
|
func WithLocalStorage(v *engine.StorageEngine) Option {
|
2020-09-28 13:22:13 +00:00
|
|
|
return func(c *cfg) {
|
|
|
|
c.localStore = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func WithContainerSource(v container.Source) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.cnrSrc = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func WithNetworkMapSource(v netmap.Source) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.netMapSrc = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func WithWorkerPool(v util.WorkerPool) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.workerPool = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func WithLocalAddressSource(v network.LocalAddressSource) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.localAddrSrc = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func WithHeadService(v *headsvc.Service) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.headSvc = v
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-12-07 17:49:47 +00:00
|
|
|
func WithRangeService(v *getsvc.Service) Option {
|
2020-09-28 13:22:13 +00:00
|
|
|
return func(c *cfg) {
|
|
|
|
c.rangeSvc = v
|
|
|
|
}
|
|
|
|
}
|
2020-11-18 13:04:59 +00:00
|
|
|
|
|
|
|
func WithClientCache(v *cache.ClientCache) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.clientCache = v
|
|
|
|
}
|
|
|
|
}
|
2020-11-23 11:24:51 +00:00
|
|
|
|
|
|
|
func WithLogger(l *logger.Logger) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.log = l
|
|
|
|
}
|
|
|
|
}
|
2020-11-23 12:59:06 +00:00
|
|
|
|
|
|
|
func WithClientOptions(opts ...client.Option) Option {
|
|
|
|
return func(c *cfg) {
|
|
|
|
c.clientOpts = opts
|
|
|
|
}
|
|
|
|
}
|