lotus/cmd/lotus-seal-worker/main.go

170 lines
3.6 KiB
Go
Raw Normal View History

package main
2019-11-21 00:52:59 +00:00
import (
"os"
2020-02-04 19:04:49 +00:00
"sync"
2019-11-21 00:52:59 +00:00
2020-02-04 19:17:18 +00:00
paramfetch "github.com/filecoin-project/go-paramfetch"
2020-02-04 19:04:49 +00:00
"github.com/filecoin-project/go-sectorbuilder"
"github.com/mitchellh/go-homedir"
logging "github.com/ipfs/go-log/v2"
2019-11-21 00:52:59 +00:00
"golang.org/x/xerrors"
"gopkg.in/urfave/cli.v2"
2020-02-04 19:17:18 +00:00
manet "github.com/multiformats/go-multiaddr-net"
"github.com/filecoin-project/lotus/api"
2019-11-21 00:52:59 +00:00
"github.com/filecoin-project/lotus/build"
lcli "github.com/filecoin-project/lotus/cli"
2020-01-08 13:49:34 +00:00
"github.com/filecoin-project/lotus/lib/lotuslog"
"github.com/filecoin-project/lotus/node/repo"
2019-11-21 00:52:59 +00:00
)
var log = logging.Logger("main")
2020-02-04 19:04:49 +00:00
const (
workers = 1 // TODO: Configurability
transfers = 1
)
2019-11-21 00:52:59 +00:00
func main() {
2020-01-08 13:49:34 +00:00
lotuslog.SetupLogLevels()
2019-11-21 00:52:59 +00:00
log.Info("Starting lotus worker")
local := []*cli.Command{
runCmd,
}
app := &cli.App{
2019-11-22 16:25:56 +00:00
Name: "lotus-seal-worker",
2019-11-21 00:52:59 +00:00
Usage: "Remote storage miner worker",
Version: build.UserVersion,
2019-11-21 00:52:59 +00:00
Flags: []cli.Flag{
&cli.StringFlag{
Name: "repo",
EnvVars: []string{"WORKER_PATH"},
Value: "~/.lotusworker", // TODO: Consider XDG_DATA_HOME
},
&cli.StringFlag{
2019-11-21 16:10:04 +00:00
Name: "storagerepo",
EnvVars: []string{"LOTUS_STORAGE_PATH"},
2019-11-21 00:52:59 +00:00
Value: "~/.lotusstorage", // TODO: Consider XDG_DATA_HOME
},
2019-12-07 14:19:46 +00:00
&cli.BoolFlag{
Name: "enable-gpu-proving",
Usage: "enable use of GPU for mining operations",
Value: true,
},
&cli.BoolFlag{
Name: "no-precommit",
},
&cli.BoolFlag{
Name: "no-commit",
},
2019-11-21 00:52:59 +00:00
},
Commands: local,
}
app.Setup()
app.Metadata["repoType"] = repo.StorageMiner
2019-11-21 00:52:59 +00:00
if err := app.Run(os.Args); err != nil {
2019-11-21 18:38:43 +00:00
log.Warnf("%+v", err)
2019-11-21 00:52:59 +00:00
return
}
}
2020-02-04 19:04:49 +00:00
type limits struct {
workLimit chan struct{}
transferLimit chan struct{}
}
2019-11-21 00:52:59 +00:00
var runCmd = &cli.Command{
Name: "run",
Usage: "Start lotus worker",
2019-11-21 00:52:59 +00:00
Action: func(cctx *cli.Context) error {
2019-12-07 14:19:46 +00:00
if !cctx.Bool("enable-gpu-proving") {
os.Setenv("BELLMAN_NO_GPU", "true")
}
2019-11-21 00:52:59 +00:00
nodeApi, closer, err := lcli.GetStorageMinerAPI(cctx)
if err != nil {
2019-11-21 18:38:43 +00:00
return xerrors.Errorf("getting miner api: %w", err)
2019-11-21 00:52:59 +00:00
}
defer closer()
ctx := lcli.ReqContext(cctx)
ainfo, err := lcli.GetAPIInfo(cctx, repo.StorageMiner)
2019-11-21 16:10:04 +00:00
if err != nil {
return xerrors.Errorf("could not get api info: %w", err)
2019-11-21 16:10:04 +00:00
}
_, storageAddr, err := manet.DialArgs(ainfo.Addr)
2019-11-21 16:10:04 +00:00
2019-11-21 18:38:43 +00:00
r, err := homedir.Expand(cctx.String("repo"))
2019-11-21 16:10:04 +00:00
if err != nil {
return err
}
2019-11-21 00:52:59 +00:00
v, err := nodeApi.Version(ctx)
if err != nil {
return err
}
if v.APIVersion != build.APIVersion {
return xerrors.Errorf("lotus-storage-miner API version doesn't match: local: ", api.Version{APIVersion: build.APIVersion})
2019-11-21 00:52:59 +00:00
}
go func() {
<-ctx.Done()
2019-12-04 16:53:32 +00:00
log.Warn("Shutting down..")
2019-11-21 00:52:59 +00:00
}()
2020-02-04 19:04:49 +00:00
limiter := &limits{
workLimit: make(chan struct{}, workers),
transferLimit: make(chan struct{}, transfers),
}
act, err := nodeApi.ActorAddress(ctx)
if err != nil {
return err
}
ssize, err := nodeApi.ActorSectorSize(ctx, act)
if err != nil {
return err
}
2020-02-04 19:17:18 +00:00
if err := paramfetch.GetParams(build.ParametersJson(), ssize); err != nil {
return xerrors.Errorf("get params: %w", err)
}
2020-02-04 19:04:49 +00:00
sb, err := sectorbuilder.NewStandalone(&sectorbuilder.Config{
SectorSize: ssize,
Miner: act,
WorkerThreads: workers,
Paths: sectorbuilder.SimplePath(r),
})
if err != nil {
return err
}
nQueues := workers + transfers
var wg sync.WaitGroup
wg.Add(nQueues)
for i := 0; i < nQueues; i++ {
go func() {
defer wg.Done()
if err := acceptJobs(ctx, nodeApi, sb, limiter, "http://"+storageAddr, ainfo.AuthHeader(), r, cctx.Bool("no-precommit"), cctx.Bool("no-commit")); err != nil {
log.Warnf("%+v", err)
return
}
}()
}
wg.Wait()
return nil
2019-11-21 00:52:59 +00:00
},
}