Merge pull request #5023 from filecoin-project/feat/worker-set-task-types
worker: Support setting task types at runtime
This commit is contained in:
@@ -2,12 +2,14 @@ package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
|
||||
"github.com/urfave/cli/v2"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"github.com/filecoin-project/lotus/chain/types"
|
||||
lcli "github.com/filecoin-project/lotus/cli"
|
||||
"github.com/filecoin-project/lotus/extern/sector-storage/sealtasks"
|
||||
)
|
||||
|
||||
var infoCmd = &cli.Command{
|
||||
@@ -49,10 +51,22 @@ var infoCmd = &cli.Command{
|
||||
return xerrors.Errorf("getting info: %w", err)
|
||||
}
|
||||
|
||||
tt, err := api.TaskTypes(ctx)
|
||||
if err != nil {
|
||||
return xerrors.Errorf("getting task types: %w", err)
|
||||
}
|
||||
|
||||
fmt.Printf("Hostname: %s\n", info.Hostname)
|
||||
fmt.Printf("CPUs: %d; GPUs: %v\n", info.Resources.CPUs, info.Resources.GPUs)
|
||||
fmt.Printf("RAM: %s; Swap: %s\n", types.SizeStr(types.NewInt(info.Resources.MemPhysical)), types.SizeStr(types.NewInt(info.Resources.MemSwap)))
|
||||
fmt.Printf("Reserved memory: %s\n", types.SizeStr(types.NewInt(info.Resources.MemReserved)))
|
||||
|
||||
fmt.Printf("Task types: ")
|
||||
for _, t := range ttList(tt) {
|
||||
fmt.Printf("%s ", t.Short())
|
||||
}
|
||||
fmt.Println()
|
||||
|
||||
fmt.Println()
|
||||
|
||||
paths, err := api.Paths(ctx)
|
||||
@@ -80,3 +94,14 @@ var infoCmd = &cli.Command{
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
func ttList(tt map[sealtasks.TaskType]struct{}) []sealtasks.TaskType {
|
||||
tasks := make([]sealtasks.TaskType, 0, len(tt))
|
||||
for taskType := range tt {
|
||||
tasks = append(tasks, taskType)
|
||||
}
|
||||
sort.Slice(tasks, func(i, j int) bool {
|
||||
return tasks[i].Less(tasks[j])
|
||||
})
|
||||
return tasks
|
||||
}
|
||||
|
||||
@@ -59,6 +59,7 @@ func main() {
|
||||
storageCmd,
|
||||
setCmd,
|
||||
waitQuietCmd,
|
||||
tasksCmd,
|
||||
}
|
||||
|
||||
app := &cli.App{
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
|
||||
"github.com/urfave/cli/v2"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"github.com/filecoin-project/lotus/api"
|
||||
lcli "github.com/filecoin-project/lotus/cli"
|
||||
"github.com/filecoin-project/lotus/extern/sector-storage/sealtasks"
|
||||
)
|
||||
|
||||
var tasksCmd = &cli.Command{
|
||||
Name: "tasks",
|
||||
Usage: "Manage task processing",
|
||||
Subcommands: []*cli.Command{
|
||||
tasksEnableCmd,
|
||||
tasksDisableCmd,
|
||||
},
|
||||
}
|
||||
|
||||
var allowSetting = map[sealtasks.TaskType]struct{}{
|
||||
sealtasks.TTAddPiece: {},
|
||||
sealtasks.TTPreCommit1: {},
|
||||
sealtasks.TTPreCommit2: {},
|
||||
sealtasks.TTCommit2: {},
|
||||
sealtasks.TTUnseal: {},
|
||||
}
|
||||
|
||||
var settableStr = func() string {
|
||||
var s []string
|
||||
for _, tt := range ttList(allowSetting) {
|
||||
s = append(s, tt.Short())
|
||||
}
|
||||
return strings.Join(s, "|")
|
||||
}()
|
||||
|
||||
var tasksEnableCmd = &cli.Command{
|
||||
Name: "enable",
|
||||
Usage: "Enable a task type",
|
||||
ArgsUsage: "[" + settableStr + "]",
|
||||
Action: taskAction(api.WorkerAPI.TaskEnable),
|
||||
}
|
||||
|
||||
var tasksDisableCmd = &cli.Command{
|
||||
Name: "disable",
|
||||
Usage: "Disable a task type",
|
||||
ArgsUsage: "[" + settableStr + "]",
|
||||
Action: taskAction(api.WorkerAPI.TaskDisable),
|
||||
}
|
||||
|
||||
func taskAction(tf func(a api.WorkerAPI, ctx context.Context, tt sealtasks.TaskType) error) func(cctx *cli.Context) error {
|
||||
return func(cctx *cli.Context) error {
|
||||
if cctx.NArg() != 1 {
|
||||
return xerrors.Errorf("expected 1 argument")
|
||||
}
|
||||
|
||||
var tt sealtasks.TaskType
|
||||
for taskType := range allowSetting {
|
||||
if taskType.Short() == cctx.Args().First() {
|
||||
tt = taskType
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if tt == "" {
|
||||
return xerrors.Errorf("unknown task type '%s'", cctx.Args().First())
|
||||
}
|
||||
|
||||
api, closer, err := lcli.GetWorkerAPI(cctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer closer()
|
||||
|
||||
ctx := lcli.ReqContext(cctx)
|
||||
|
||||
return tf(api, ctx, tt)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user