49 lines
1.3 KiB
Go
49 lines
1.3 KiB
Go
package sectorstorage
|
|
|
|
import (
|
|
"context"
|
|
|
|
"golang.org/x/xerrors"
|
|
|
|
"github.com/filecoin-project/specs-actors/actors/abi"
|
|
|
|
"github.com/filecoin-project/lotus/extern/sector-storage/sealtasks"
|
|
"github.com/filecoin-project/lotus/extern/sector-storage/stores"
|
|
)
|
|
|
|
type taskSelector struct {
|
|
best []stores.StorageInfo //nolint: unused, structcheck
|
|
}
|
|
|
|
func newTaskSelector() *taskSelector {
|
|
return &taskSelector{}
|
|
}
|
|
|
|
func (s *taskSelector) Ok(ctx context.Context, task sealtasks.TaskType, spt abi.RegisteredSealProof, whnd *workerHandle) (bool, error) {
|
|
tasks, err := whnd.w.TaskTypes(ctx)
|
|
if err != nil {
|
|
return false, xerrors.Errorf("getting supported worker task types: %w", err)
|
|
}
|
|
_, supported := tasks[task]
|
|
|
|
return supported, nil
|
|
}
|
|
|
|
func (s *taskSelector) Cmp(ctx context.Context, _ sealtasks.TaskType, a, b *workerHandle) (bool, error) {
|
|
atasks, err := a.w.TaskTypes(ctx)
|
|
if err != nil {
|
|
return false, xerrors.Errorf("getting supported worker task types: %w", err)
|
|
}
|
|
btasks, err := b.w.TaskTypes(ctx)
|
|
if err != nil {
|
|
return false, xerrors.Errorf("getting supported worker task types: %w", err)
|
|
}
|
|
if len(atasks) != len(btasks) {
|
|
return len(atasks) < len(btasks), nil // prefer workers which can do less
|
|
}
|
|
|
|
return a.utilization() < b.utilization(), nil
|
|
}
|
|
|
|
var _ WorkerSelector = &allocSelector{}
|