sectorbuilder: check free space before creating sectors
This commit is contained in:
@@ -17,11 +17,11 @@ func (sb *SectorBuilder) SectorName(sectorID uint64) string {
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) StagedSectorPath(sectorID uint64) string {
|
||||
return filepath.Join(sb.stagedDir, sb.SectorName(sectorID))
|
||||
return filepath.Join(sb.filesystem.staging(), sb.SectorName(sectorID))
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) unsealedSectorPath(sectorID uint64) string {
|
||||
return filepath.Join(sb.unsealedDir, sb.SectorName(sectorID))
|
||||
return filepath.Join(sb.filesystem.unsealed(), sb.SectorName(sectorID))
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) stagedSectorFile(sectorID uint64) (*os.File, error) {
|
||||
@@ -29,13 +29,13 @@ func (sb *SectorBuilder) stagedSectorFile(sectorID uint64) (*os.File, error) {
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) SealedSectorPath(sectorID uint64) (string, error) {
|
||||
path := filepath.Join(sb.sealedDir, sb.SectorName(sectorID))
|
||||
path := filepath.Join(sb.filesystem.sealed(), sb.SectorName(sectorID))
|
||||
|
||||
return path, nil
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) sectorCacheDir(sectorID uint64) (string, error) {
|
||||
dir := filepath.Join(sb.cacheDir, sb.SectorName(sectorID))
|
||||
dir := filepath.Join(sb.filesystem.cache(), sb.SectorName(sectorID))
|
||||
|
||||
err := os.Mkdir(dir, 0755)
|
||||
if os.IsExist(err) {
|
||||
@@ -48,11 +48,11 @@ func (sb *SectorBuilder) sectorCacheDir(sectorID uint64) (string, error) {
|
||||
func (sb *SectorBuilder) GetPath(typ string, sectorName string) (string, error) {
|
||||
switch typ {
|
||||
case "staged":
|
||||
return filepath.Join(sb.stagedDir, sectorName), nil
|
||||
return filepath.Join(sb.filesystem.staging(), sectorName), nil
|
||||
case "sealed":
|
||||
return filepath.Join(sb.sealedDir, sectorName), nil
|
||||
return filepath.Join(sb.filesystem.sealed(), sectorName), nil
|
||||
case "cache":
|
||||
return filepath.Join(sb.cacheDir, sectorName), nil
|
||||
return filepath.Join(sb.filesystem.cache(), sectorName), nil
|
||||
default:
|
||||
return "", xerrors.Errorf("unknown sector type for write: %s", typ)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
package sectorbuilder
|
||||
|
||||
import (
|
||||
"github.com/filecoin-project/lotus/chain/types"
|
||||
"golang.org/x/xerrors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"syscall"
|
||||
)
|
||||
|
||||
type dataType int
|
||||
|
||||
const (
|
||||
dataCache dataType = iota
|
||||
dataStaging
|
||||
dataSealed
|
||||
dataUnsealed
|
||||
|
||||
nDataTypes
|
||||
)
|
||||
|
||||
var overheadMul = []uint64{ // * sectorSize
|
||||
dataCache: 11, // TODO: check if true for 32G sectors
|
||||
dataStaging: 1,
|
||||
dataSealed: 1,
|
||||
dataUnsealed: 1,
|
||||
}
|
||||
|
||||
type fs struct {
|
||||
path string
|
||||
|
||||
// in progress actions
|
||||
|
||||
reserved [nDataTypes]uint64
|
||||
|
||||
lk sync.Mutex
|
||||
}
|
||||
|
||||
func openFs(dir string) *fs {
|
||||
return &fs{
|
||||
path: dir,
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fs) init() error {
|
||||
for _, dir := range []string{f.path, f.cache(), f.staging(), f.sealed(), f.unsealed()} {
|
||||
if err := os.Mkdir(dir, 0755); err != nil {
|
||||
if os.IsExist(err) {
|
||||
continue
|
||||
}
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fs) cache() string {
|
||||
return filepath.Join(f.path, "cache")
|
||||
}
|
||||
|
||||
func (f *fs) staging() string {
|
||||
return filepath.Join(f.path, "staging")
|
||||
}
|
||||
|
||||
func (f *fs) sealed() string {
|
||||
return filepath.Join(f.path, "sealed")
|
||||
}
|
||||
|
||||
func (f *fs) unsealed() string {
|
||||
return filepath.Join(f.path, "unsealed")
|
||||
}
|
||||
|
||||
func (f *fs) reservedBytes() int64 {
|
||||
var out int64
|
||||
for _, r := range f.reserved {
|
||||
out += int64(r)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func (f *fs) reserve(typ dataType, size uint64) error {
|
||||
f.lk.Lock()
|
||||
defer f.lk.Unlock()
|
||||
|
||||
var fsstat syscall.Statfs_t
|
||||
|
||||
if err := syscall.Statfs(f.path, &fsstat); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
avail := int64(fsstat.Bavail) * fsstat.Bsize
|
||||
|
||||
avail -= f.reservedBytes()
|
||||
|
||||
need := overheadMul[typ] * size
|
||||
|
||||
if int64(need) > avail {
|
||||
return xerrors.Errorf("not enough space in '%s', need %s, available %s", f.path, types.NewInt(need).SizeStr(), types.NewInt(uint64(avail)).SizeStr())
|
||||
}
|
||||
|
||||
f.reserved[typ] += need
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fs) free(typ dataType, sectorSize uint64) {
|
||||
f.lk.Lock()
|
||||
defer f.lk.Unlock()
|
||||
|
||||
f.reserved[typ] -= overheadMul[typ] * sectorSize
|
||||
|
||||
return
|
||||
}
|
||||
@@ -1,8 +1,6 @@
|
||||
package sectorbuilder
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
|
||||
"github.com/filecoin-project/lotus/chain/address"
|
||||
"github.com/filecoin-project/lotus/node/modules/dtypes"
|
||||
)
|
||||
@@ -13,18 +11,10 @@ func TempSectorbuilderDir(dir string, sectorSize uint64, ds dtypes.MetadataDS) (
|
||||
return nil, err
|
||||
}
|
||||
|
||||
unsealed := filepath.Join(dir, "unsealed")
|
||||
sealed := filepath.Join(dir, "sealed")
|
||||
staging := filepath.Join(dir, "staging")
|
||||
cache := filepath.Join(dir, "cache")
|
||||
|
||||
sb, err := New(&Config{
|
||||
SectorSize: sectorSize,
|
||||
|
||||
SealedDir: sealed,
|
||||
StagedDir: staging,
|
||||
UnsealedDir: unsealed,
|
||||
CacheDir: cache,
|
||||
Dir: dir,
|
||||
|
||||
WorkerThreads: 2,
|
||||
Miner: addr,
|
||||
|
||||
@@ -64,11 +64,6 @@ type SectorBuilder struct {
|
||||
|
||||
Miner address.Address
|
||||
|
||||
stagedDir string
|
||||
sealedDir string
|
||||
cacheDir string
|
||||
unsealedDir string
|
||||
|
||||
unsealLk sync.Mutex
|
||||
|
||||
noCommit bool
|
||||
@@ -89,6 +84,9 @@ type SectorBuilder struct {
|
||||
commitWait int32
|
||||
unsealWait int32
|
||||
|
||||
fsLk sync.Mutex
|
||||
filesystem *fs // TODO: multi-fs support
|
||||
|
||||
stopping chan struct{}
|
||||
}
|
||||
|
||||
@@ -135,11 +133,8 @@ type Config struct {
|
||||
NoCommit bool
|
||||
NoPreCommit bool
|
||||
|
||||
CacheDir string
|
||||
SealedDir string
|
||||
StagedDir string
|
||||
UnsealedDir string
|
||||
_ struct{} // guard against nameless init
|
||||
Dir string
|
||||
_ struct{} // guard against nameless init
|
||||
}
|
||||
|
||||
func New(cfg *Config, ds dtypes.MetadataDS) (*SectorBuilder, error) {
|
||||
@@ -147,15 +142,6 @@ func New(cfg *Config, ds dtypes.MetadataDS) (*SectorBuilder, error) {
|
||||
return nil, xerrors.Errorf("minimum worker threads is %d, specified %d", PoStReservedWorkers, cfg.WorkerThreads)
|
||||
}
|
||||
|
||||
for _, dir := range []string{cfg.StagedDir, cfg.SealedDir, cfg.CacheDir, cfg.UnsealedDir} {
|
||||
if err := os.Mkdir(dir, 0755); err != nil {
|
||||
if os.IsExist(err) {
|
||||
continue
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
var lastUsedID uint64
|
||||
b, err := ds.Get(lastSectorIdKey)
|
||||
switch err {
|
||||
@@ -185,10 +171,7 @@ func New(cfg *Config, ds dtypes.MetadataDS) (*SectorBuilder, error) {
|
||||
ssize: cfg.SectorSize,
|
||||
lastID: lastUsedID,
|
||||
|
||||
stagedDir: cfg.StagedDir,
|
||||
sealedDir: cfg.SealedDir,
|
||||
cacheDir: cfg.CacheDir,
|
||||
unsealedDir: cfg.UnsealedDir,
|
||||
filesystem: openFs(cfg.Dir),
|
||||
|
||||
Miner: cfg.Miner,
|
||||
|
||||
@@ -205,35 +188,33 @@ func New(cfg *Config, ds dtypes.MetadataDS) (*SectorBuilder, error) {
|
||||
stopping: make(chan struct{}),
|
||||
}
|
||||
|
||||
if err := sb.filesystem.init(); err != nil {
|
||||
return nil, xerrors.Errorf("initializing sectorbuilder filesystem: %w", err)
|
||||
}
|
||||
|
||||
return sb, nil
|
||||
}
|
||||
|
||||
func NewStandalone(cfg *Config) (*SectorBuilder, error) {
|
||||
for _, dir := range []string{cfg.StagedDir, cfg.SealedDir, cfg.CacheDir, cfg.UnsealedDir} {
|
||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||
if os.IsExist(err) {
|
||||
continue
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
return &SectorBuilder{
|
||||
sb := &SectorBuilder{
|
||||
ds: nil,
|
||||
|
||||
ssize: cfg.SectorSize,
|
||||
|
||||
Miner: cfg.Miner,
|
||||
stagedDir: cfg.StagedDir,
|
||||
sealedDir: cfg.SealedDir,
|
||||
cacheDir: cfg.CacheDir,
|
||||
unsealedDir: cfg.UnsealedDir,
|
||||
Miner: cfg.Miner,
|
||||
filesystem: openFs(cfg.Dir),
|
||||
|
||||
taskCtr: 1,
|
||||
remotes: map[int]*remote{},
|
||||
rateLimit: make(chan struct{}, cfg.WorkerThreads),
|
||||
stopping: make(chan struct{}),
|
||||
}, nil
|
||||
}
|
||||
|
||||
if err := sb.filesystem.init(); err != nil {
|
||||
return nil, xerrors.Errorf("initializing sectorbuilder filesystem: %w", err)
|
||||
}
|
||||
|
||||
return sb, nil
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) checkRateLimit() {
|
||||
@@ -312,6 +293,13 @@ func (sb *SectorBuilder) AcquireSectorId() (uint64, error) {
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) AddPiece(pieceSize uint64, sectorId uint64, file io.Reader, existingPieceSizes []uint64) (PublicPieceInfo, error) {
|
||||
fs := sb.filesystem
|
||||
|
||||
if err := fs.reserve(dataStaging, sb.ssize); err != nil {
|
||||
return PublicPieceInfo{}, err
|
||||
}
|
||||
defer fs.free(dataStaging, sb.ssize)
|
||||
|
||||
atomic.AddInt32(&sb.addPieceWait, 1)
|
||||
ret := sb.RateLimit()
|
||||
atomic.AddInt32(&sb.addPieceWait, -1)
|
||||
@@ -347,6 +335,13 @@ func (sb *SectorBuilder) AddPiece(pieceSize uint64, sectorId uint64, file io.Rea
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) ReadPieceFromSealedSector(sectorID uint64, offset uint64, size uint64, ticket []byte, commD []byte) (io.ReadCloser, error) {
|
||||
fs := sb.filesystem
|
||||
|
||||
if err := fs.reserve(dataUnsealed, sb.ssize); err != nil { // TODO: this needs to get smarter when we start supporting partial unseals
|
||||
return nil, err
|
||||
}
|
||||
defer fs.free(dataUnsealed, sb.ssize)
|
||||
|
||||
atomic.AddInt32(&sb.unsealWait, 1)
|
||||
// TODO: Don't wait if cached
|
||||
ret := sb.RateLimit() // TODO: check perf, consider remote unseal worker
|
||||
@@ -433,6 +428,18 @@ func (sb *SectorBuilder) sealPreCommitRemote(call workerCall) (RawSealPreCommitO
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) SealPreCommit(sectorID uint64, ticket SealTicket, pieces []PublicPieceInfo) (RawSealPreCommitOutput, error) {
|
||||
fs := sb.filesystem
|
||||
|
||||
if err := fs.reserve(dataCache, sb.ssize); err != nil {
|
||||
return RawSealPreCommitOutput{}, err
|
||||
}
|
||||
defer fs.free(dataCache, sb.ssize)
|
||||
|
||||
if err := fs.reserve(dataSealed, sb.ssize); err != nil {
|
||||
return RawSealPreCommitOutput{}, err
|
||||
}
|
||||
defer fs.free(dataSealed, sb.ssize)
|
||||
|
||||
call := workerCall{
|
||||
task: WorkerTask{
|
||||
Type: WorkerPreCommit,
|
||||
@@ -692,15 +699,15 @@ func fallbackPostChallengeCount(sectors uint64) uint64 {
|
||||
}
|
||||
|
||||
func (sb *SectorBuilder) ImportFrom(osb *SectorBuilder, symlink bool) error {
|
||||
if err := migrate(osb.cacheDir, sb.cacheDir, symlink); err != nil {
|
||||
if err := migrate(osb.filesystem.cache(), sb.filesystem.cache(), symlink); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := migrate(osb.sealedDir, sb.sealedDir, symlink); err != nil {
|
||||
if err := migrate(osb.filesystem.sealed(), sb.filesystem.sealed(), symlink); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := migrate(osb.stagedDir, sb.stagedDir, symlink); err != nil {
|
||||
if err := migrate(osb.filesystem.staging(), sb.filesystem.staging(), symlink); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user