diff --git a/cmd/image/build.go b/cmd/image/build.go index a86ee73..4dfff6c 100644 --- a/cmd/image/build.go +++ b/cmd/image/build.go @@ -33,7 +33,6 @@ func (o *buildImageOptions) run(ctx context.Context) error { interactive := term.IsTerminal(int(os.Stdout.Fd())) progress := map[string]int{} - p := 0 for { msg, err := resp.Recv() if errors.Is(err, io.EOF) { @@ -50,7 +49,6 @@ func (o *buildImageOptions) run(ctx context.Context) error { fmt.Print(msg.Stream) if msg.Status == "finished" { clear(progress) - p = 0 } case msg.Status != "": if msg.Id == "" { @@ -58,11 +56,10 @@ func (o *buildImageOptions) run(ctx context.Context) error { } else { data := fmt.Sprintf("%s: %s %s", msg.Id, msg.Status, msg.Progress) if pos, ok := progress[msg.Id]; !ok { - progress[msg.Id] = p + progress[msg.Id] = len(progress) fmt.Println(data) - p++ } else if interactive { - fmt.Printf(progressRewrite, p-pos, data) + fmt.Printf(progressRewrite, len(progress)-pos, data) } else { fmt.Println(data) } diff --git a/cmd/image/cache.go b/cmd/image/cache.go index c91a378..1a015c2 100644 --- a/cmd/image/cache.go +++ b/cmd/image/cache.go @@ -13,20 +13,13 @@ import ( ) type cacheImageOptions struct { - client corepb.CoreRPCClient - images []string - podname string - nodenames []string + client corepb.CoreRPCClient + opts *corepb.CacheImageOptions } func (o *cacheImageOptions) run(ctx context.Context) error { logger := log.WithFunc("image.cacheImageOptions.run") - opts := &corepb.CacheImageOptions{ - Images: o.images, - Podname: o.podname, - Nodenames: o.nodenames, - } - resp, err := o.client.CacheImage(ctx, opts) + resp, err := o.client.CacheImage(ctx, o.opts) if err != nil { return err } @@ -52,10 +45,12 @@ func cmdImageCache(ctx context.Context, cmd *cli.Command) error { } o := &cacheImageOptions{ - client: client, - images: images, - podname: cmd.String(flagPod), - nodenames: cmd.StringSlice(flagNode), + client: client, + opts: &corepb.CacheImageOptions{ + Images: images, + Podname: cmd.String(flagPod), + Nodenames: cmd.StringSlice(flagNode), + }, } return o.run(ctx) } diff --git a/cmd/image/remove.go b/cmd/image/remove.go index ee459dc..0ecf82d 100644 --- a/cmd/image/remove.go +++ b/cmd/image/remove.go @@ -13,22 +13,13 @@ import ( ) type removeImageOptions struct { - client corepb.CoreRPCClient - images []string - podname string - nodenames []string - prune bool + client corepb.CoreRPCClient + opts *corepb.RemoveImageOptions } func (o *removeImageOptions) run(ctx context.Context) error { logger := log.WithFunc("image.removeImageOptions.run") - opts := &corepb.RemoveImageOptions{ - Images: o.images, - Podname: o.podname, - Nodenames: o.nodenames, - Prune: o.prune, - } - resp, err := o.client.RemoveImage(ctx, opts) + resp, err := o.client.RemoveImage(ctx, o.opts) if err != nil { return err } @@ -57,11 +48,13 @@ func cmdImageRemove(ctx context.Context, cmd *cli.Command) error { } o := &removeImageOptions{ - client: client, - images: images, - podname: cmd.String(flagPod), - nodenames: cmd.StringSlice(flagNode), - prune: cmd.Bool("prune"), + client: client, + opts: &corepb.RemoveImageOptions{ + Images: images, + Podname: cmd.String(flagPod), + Nodenames: cmd.StringSlice(flagNode), + Prune: cmd.Bool("prune"), + }, } return o.run(ctx) } diff --git a/cmd/image/stream_test.go b/cmd/image/stream_test.go index 91823bd..46aa0b4 100644 --- a/cmd/image/stream_test.go +++ b/cmd/image/stream_test.go @@ -114,7 +114,7 @@ func TestCacheImageReportsFailureReason(t *testing.T) { client: &fakeImageClient{cache: &fakeStream[corepb.CacheImageMessage]{msgs: []*corepb.CacheImageMessage{ {Image: "app:v1", Nodename: "node1", Success: false, Message: "no such image"}, }}}, - images: []string{"app:v1"}, + opts: &corepb.CacheImageOptions{Images: []string{"app:v1"}}, } err := o.run(t.Context()) @@ -128,7 +128,7 @@ func TestRemoveImageReportsFailure(t *testing.T) { client: &fakeImageClient{remove: &fakeStream[corepb.RemoveImageMessage]{msgs: []*corepb.RemoveImageMessage{ {Image: "app:v1", Success: false}, }}}, - images: []string{"app:v1"}, + opts: &corepb.RemoveImageOptions{Images: []string{"app:v1"}}, } err := o.run(t.Context()) @@ -168,14 +168,14 @@ func TestImageCommandsPassANodeOnlyRequestThrough(t *testing.T) { { name: "cache", run: func(ctx context.Context, c *fakeImageClient) error { - o := &cacheImageOptions{client: c, images: []string{"app:v1"}, nodenames: []string{"node1"}} + o := &cacheImageOptions{client: c, opts: &corepb.CacheImageOptions{Images: []string{"app:v1"}, Nodenames: []string{"node1"}}} return o.run(ctx) }, }, { name: "remove", run: func(ctx context.Context, c *fakeImageClient) error { - o := &removeImageOptions{client: c, images: []string{"app:v1"}, nodenames: []string{"node1"}} + o := &removeImageOptions{client: c, opts: &corepb.RemoveImageOptions{Images: []string{"app:v1"}, Nodenames: []string{"node1"}}} return o.run(ctx) }, }, diff --git a/cmd/lambda/run.go b/cmd/lambda/run.go index 4ff80c4..7d6d4ba 100644 --- a/cmd/lambda/run.go +++ b/cmd/lambda/run.go @@ -19,7 +19,6 @@ var newline = []byte{'\n'} type runLambdaOptions struct { client corepb.CoreRPCClient opts *corepb.RunAndWaitOptions - stdin bool printWorkloadID bool } @@ -52,7 +51,7 @@ func (o *runLambdaOptions) lambda(ctx context.Context) (int, error) { _ = iStream.Send(newline) }() - exitCount, stdin := int(o.opts.GetDeployOptions().GetCount()), o.stdin + exitCount, stdin := int(o.opts.GetDeployOptions().GetCount()), o.opts.GetDeployOptions().GetOpenStdin() if o.opts.Async { exitCount, stdin = 0, false } @@ -73,7 +72,6 @@ func cmdLambdaRun(ctx context.Context, cmd *cli.Command) error { o := &runLambdaOptions{ client: client, opts: opts, - stdin: cmd.Bool("stdin"), printWorkloadID: cmd.Bool("workload-id"), } return o.run(ctx) @@ -86,11 +84,11 @@ func generateLambdaOptions(cmd *cli.Command) (*corepb.RunAndWaitOptions, error) network := cmd.String("network") - memoryRequest, err := utils.ParseRAMInHuman(cmd.String("memory-request")) + memoryRequest, err := resourcetypes.ParseRAMInHuman(cmd.String("memory-request")) if err != nil { return nil, fmt.Errorf("parse memory-request: %w", err) } - memoryLimit, err := utils.ParseRAMInHuman(cmd.String("memory")) + memoryLimit, err := resourcetypes.ParseRAMInHuman(cmd.String("memory")) if err != nil { return nil, fmt.Errorf("parse memory: %w", err) } @@ -111,11 +109,11 @@ func generateLambdaOptions(cmd *cli.Command) (*corepb.RunAndWaitOptions, error) "memory-request": memoryRequest, "memory-limit": memoryLimit, } - storageRequest, err := utils.ParseRAMInHuman(cmd.String("storage-request")) + storageRequest, err := resourcetypes.ParseRAMInHuman(cmd.String("storage-request")) if err != nil { return nil, fmt.Errorf("parse storage-request: %w", err) } - storageLimit, err := utils.ParseRAMInHuman(cmd.String("storage")) + storageLimit, err := resourcetypes.ParseRAMInHuman(cmd.String("storage")) if err != nil { return nil, fmt.Errorf("parse storage: %w", err) } diff --git a/cmd/node/status.go b/cmd/node/status.go index 6414d51..d3fd601 100644 --- a/cmd/node/status.go +++ b/cmd/node/status.go @@ -21,11 +21,13 @@ type setNodeStatusOptions struct { } func (o *setNodeStatusOptions) run(ctx context.Context) error { + err := o.heartbeat(ctx) if o.interval == 0 { - return o.heartbeat(ctx) + return err } logger := log.WithFunc("node.setNodeStatusOptions.run") + logger.Error(ctx, err, "heartbeat") ticker := time.NewTicker(time.Duration(o.interval) * time.Second) defer ticker.Stop() diff --git a/cmd/pod/capacity.go b/cmd/pod/capacity.go index bd252ca..121942e 100644 --- a/cmd/pod/capacity.go +++ b/cmd/pod/capacity.go @@ -70,12 +70,12 @@ func cmdPodCapacity(ctx context.Context, cmd *cli.Command) error { } func capacityResources(cmd *cli.Command) (map[string][]byte, error) { - memory, err := utils.ParseRAMInHuman(cmd.String(flagMemory)) + memory, err := resourcetypes.ParseRAMInHuman(cmd.String(flagMemory)) if err != nil { return nil, fmt.Errorf("parse memory: %w", err) } - storage, err := utils.ParseRAMInHuman(cmd.String(flagStorage)) + storage, err := resourcetypes.ParseRAMInHuman(cmd.String(flagStorage)) if err != nil { return nil, fmt.Errorf("parse storage: %w", err) } diff --git a/cmd/pod/cmd.go b/cmd/pod/cmd.go index 4a5cc98..4f3828f 100644 --- a/cmd/pod/cmd.go +++ b/cmd/pod/cmd.go @@ -38,7 +38,6 @@ func Command() *cli.Command { &cli.StringFlag{ Name: "desc", Usage: "description of pod", - Value: "", }, }, }, @@ -92,13 +91,11 @@ func Command() *cli.Command { &cli.BoolFlag{ Name: "cpu-bind", Usage: "bind cpu or not", - Value: false, }, &cli.StringSliceFlag{ - Name: "node", - Aliases: []string{"n"}, - Usage: "Specified the node(s) should join into the calculation. Could be specified multiple times with different names", - Required: false, + Name: "node", + Aliases: []string{"n"}, + Usage: "Specified the node(s) should join into the calculation. Could be specified multiple times with different names", }, utils.ExtraResourcesFlag(), }, diff --git a/cmd/pod/nodes.go b/cmd/pod/nodes.go index ce41d56..6b33b0f 100644 --- a/cmd/pod/nodes.go +++ b/cmd/pod/nodes.go @@ -48,6 +48,11 @@ func cmdPodListNodes(ctx context.Context, cmd *cli.Command) error { return err } + name := cmd.Args().First() + if name == "" { + return errors.New("pod name must be given") + } + filter := strings.ToLower(cmd.String("filter")) if filter != up && filter != down && filter != all { return errors.New("filter should be one of up/down/all") @@ -55,7 +60,7 @@ func cmdPodListNodes(ctx context.Context, cmd *cli.Command) error { o := &listPodNodesOptions{ client: client, - name: cmd.Args().First(), + name: name, filter: filter, labels: utils.SplitEquality(cmd.StringSlice("label")), timeoutInSecond: int32(cmd.Int("timeout")), //nolint:gosec diff --git a/cmd/pod/resource.go b/cmd/pod/resource.go index 743c39a..8db5406 100644 --- a/cmd/pod/resource.go +++ b/cmd/pod/resource.go @@ -8,7 +8,6 @@ import ( "strconv" "strings" - "github.com/projecteru2/core/log" corepb "github.com/projecteru2/core/rpc/gen" "github.com/urfave/cli/v3" @@ -21,56 +20,10 @@ var filterExpr = regexp.MustCompile(`^\s*(?Pcpu|memory|storage|volume)\s*( type resourcePodOptions struct { client corepb.CoreRPCClient name string - expr string + keep describe.NodeResourceFilter stream bool } -func (o *resourcePodOptions) filter(ctx context.Context, ch <-chan *corepb.NodeResource) (<-chan *corepb.NodeResource, error) { - if o.expr == "" { - return ch, nil - } - - filter := match(o.expr) - if len(filter) == 0 { - return nil, fmt.Errorf("invalid filter %q, want one of cpu/memory/storage/volume with an operator and a value", o.expr) - } - - var ( - value = filter["value"] - percent bool - ) - if v, ok := strings.CutSuffix(value, "%"); ok { - value = v - percent = true - } - - v, err := strconv.ParseFloat(value, 64) - if err != nil { - return nil, err - } - if percent { - v /= 100 - } - - rv := make(chan *corepb.NodeResource) - go func() { - defer close(rv) - logger := log.WithFunc("pod.resourcePodOptions.filter") - for nr := range ch { - l, err := attr(nr, filter["name"]) - if err != nil { - logger.Errorf(ctx, err, "resource percent of node %s", nr.Name) - continue - } - if !compare(filter["op"], l, v) { - continue - } - rv <- nr - } - }() - return rv, nil -} - func (o *resourcePodOptions) run(ctx context.Context) error { resp, err := o.client.GetPodResource(ctx, &corepb.GetPodOptions{ Name: o.name, @@ -80,12 +33,7 @@ func (o *resourcePodOptions) run(ctx context.Context) error { } ch, wait := utils.StreamToChan(resp.Recv) - resChan, err := o.filter(ctx, ch) - if err != nil { - return err - } - - describe.NodeResources(ctx, resChan, o.stream) + describe.NodeResources(ctx, ch, o.stream, o.keep) return wait() } @@ -100,15 +48,45 @@ func cmdPodResource(ctx context.Context, cmd *cli.Command) error { return errors.New("pod name must be given") } + keep, err := parseFilter(cmd.String("filter")) + if err != nil { + return err + } + o := &resourcePodOptions{ client: client, name: name, - expr: cmd.String("filter"), + keep: keep, stream: cmd.Bool("stream"), } return o.run(ctx) } +func parseFilter(expr string) (describe.NodeResourceFilter, error) { + if expr == "" { + return nil, nil + } + + filter := match(expr) + if len(filter) == 0 { + return nil, fmt.Errorf("invalid filter %q, want one of cpu/memory/storage/volume with an operator and a value", expr) + } + + value, percent := strings.CutSuffix(filter["value"], "%") + v, err := strconv.ParseFloat(value, 64) + if err != nil { + return nil, err + } + if percent { + v /= 100 + } + + name, op := filter["name"], filter["op"] + return func(cpumem, storage map[string]float64) bool { + return compare(op, attr(cpumem, storage, name), v) + }, nil +} + func match(s string) map[string]string { rv := make(map[string]string) founds := filterExpr.FindStringSubmatch(s) @@ -137,21 +115,17 @@ func compare(operator string, left, right float64) bool { } } -func attr(nr *corepb.NodeResource, name string) (float64, error) { - cr, sr, err := describe.ToResourcePercent(nr) - if err != nil { - return 0, err - } +func attr(cpumem, storage map[string]float64, name string) float64 { switch name { case flagCPU: - return cr[flagCPU], nil + return cpumem[flagCPU] case flagMemory: - return cr[flagMemory], nil + return cpumem[flagMemory] case flagStorage: - return sr[flagStorage], nil + return storage[flagStorage] case "volume": - return sr["volumes"], nil + return storage["volumes"] default: - return 0, nil + return 0 } } diff --git a/cmd/pod/resource_test.go b/cmd/pod/resource_test.go index e54bf30..5c4f0e6 100644 --- a/cmd/pod/resource_test.go +++ b/cmd/pod/resource_test.go @@ -5,7 +5,6 @@ import ( "slices" "testing" - corepb "github.com/projecteru2/core/rpc/gen" "github.com/urfave/cli/v3" ) @@ -82,10 +81,8 @@ func TestCompare(t *testing.T) { } func TestAttr(t *testing.T) { - nr := &corepb.NodeResource{ - ResourceUsage: `{"cpumem":{"cpu":2,"memory":512},"resource-storage":{"storage":50,"volumes":{"/data":10}}}`, - ResourceCapacity: `{"cpumem":{"cpu":8,"memory":2048},"resource-storage":{"storage":200,"volumes":{"/data":40}}}`, - } + cpumem := map[string]float64{flagCPU: 0.25, flagMemory: 0.25} + storage := map[string]float64{flagStorage: 0.25, "volumes": 0.25} tests := []struct { name string @@ -101,11 +98,7 @@ func TestAttr(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got, err := attr(nr, tt.attr) - if err != nil { - t.Fatalf("attr: %v", err) - } - if got != tt.want { + if got := attr(cpumem, storage, tt.attr); got != tt.want { t.Errorf("got %v, want %v", got, tt.want) } }) diff --git a/cmd/pod/stream_test.go b/cmd/pod/stream_test.go index 7118f8e..0817cc7 100644 --- a/cmd/pod/stream_test.go +++ b/cmd/pod/stream_test.go @@ -5,7 +5,6 @@ import ( "errors" "io" "os" - "slices" "testing" corepb "github.com/projecteru2/core/rpc/gen" @@ -70,8 +69,7 @@ func TestPodResourceRejectsUnparsableFilter(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - o := &resourcePodOptions{expr: tt.expr} - _, err := o.filter(t.Context(), make(chan *corepb.NodeResource)) + _, err := parseFilter(tt.expr) if (err != nil) != tt.wantErr { t.Errorf("got %v, wantErr %v", err, tt.wantErr) } @@ -79,26 +77,6 @@ func TestPodResourceRejectsUnparsableFilter(t *testing.T) { } } -func TestPodResourceFilterDropsUnparsableNodes(t *testing.T) { - o := &resourcePodOptions{expr: "cpu<=100%"} - ch := make(chan *corepb.NodeResource, 2) - ch <- &corepb.NodeResource{Name: "bad", ResourceUsage: "{"} - ch <- &corepb.NodeResource{Name: "good", ResourceUsage: "{}", ResourceCapacity: "{}"} - close(ch) - - out, err := o.filter(t.Context(), ch) - if err != nil { - t.Fatalf("filter: %v", err) - } - names := []string{} - for nr := range out { - names = append(names, nr.Name) - } - if !slices.Equal(names, []string{"good"}) { - t.Errorf("got %v, want only the parsable node", names) - } -} - func TestListDownNodesUsesOneCall(t *testing.T) { client := &fakePodClient{nodes: func() *fakeStream[corepb.Node] { return &fakeStream[corepb.Node]{msgs: []*corepb.Node{{Name: "n1", Available: true}}} diff --git a/cmd/status/status.go b/cmd/status/status.go index 50ef043..5a83ba6 100644 --- a/cmd/status/status.go +++ b/cmd/status/status.go @@ -58,7 +58,7 @@ func (o *statusOptions) run(ctx context.Context) error { case !msg.Status.Healthy: logger.Warnf(ctx, "[%s] %s on %s is unhealthy", coreutils.ShortID(msg.Id), msg.Workload.Name, msg.Workload.Nodename) default: - logger.Infof(ctx, "[%s] %s back to life", coreutils.ShortID(msg.Workload.Id), msg.Workload.Name) + logger.Infof(ctx, "[%s] %s back to life", coreutils.ShortID(msg.Id), msg.Workload.Name) for networkName, addrs := range msg.Workload.Publish { logger.Infof(ctx, "[%s] published at %s bind %v", coreutils.ShortID(msg.Id), networkName, addrs) } diff --git a/cmd/utils/file.go b/cmd/utils/file.go index 00e5a4d..b0554dc 100644 --- a/cmd/utils/file.go +++ b/cmd/utils/file.go @@ -96,20 +96,26 @@ func ReadAllFiles(files []string) (map[string]*types.LinuxFile, error) { } // SplitFiles turns a list of src:dst strings into a map. -func SplitFiles(files []string) map[string]string { +func SplitFiles(files []string) (map[string]string, error) { ret := map[string]string{} for _, f := range files { - ps := strings.Split(f, ":") - if len(ps) < 2 { - continue + src, dst, ok := strings.Cut(f, ":") + if !ok || strings.Contains(dst, ":") { + return nil, fmt.Errorf("invalid file %q, want src:dst", f) } - ret[ps[0]] = ps[1] + ret[src] = dst } - return ret + return ret, nil } -// GetSpecFromRemote fetches a spec over HTTP. -func GetSpecFromRemote(ctx context.Context, uri string) ([]byte, error) { +func ReadSpecURI(ctx context.Context, uri string) ([]byte, error) { + if strings.HasPrefix(uri, "http://") || strings.HasPrefix(uri, "https://") { + return getSpecFromRemote(ctx, uri) + } + return os.ReadFile(uri) //nolint:gosec +} + +func getSpecFromRemote(ctx context.Context, uri string) ([]byte, error) { req, err := http.NewRequestWithContext(ctx, http.MethodGet, uri, nil) if err != nil { return nil, err @@ -124,10 +130,3 @@ func GetSpecFromRemote(ctx context.Context, uri string) ([]byte, error) { } return io.ReadAll(resp.Body) } - -func ReadSpecURI(ctx context.Context, uri string) ([]byte, error) { - if strings.HasPrefix(uri, "http://") || strings.HasPrefix(uri, "https://") { - return GetSpecFromRemote(ctx, uri) - } - return os.ReadFile(uri) //nolint:gosec -} diff --git a/cmd/utils/file_test.go b/cmd/utils/file_test.go index c861677..236b0bd 100644 --- a/cmd/utils/file_test.go +++ b/cmd/utils/file_test.go @@ -67,18 +67,26 @@ func TestReadAllFiles(t *testing.T) { func TestSplitFiles(t *testing.T) { tests := []struct { - name string - files []string - want map[string]string + name string + files []string + want map[string]string + wantErr bool }{ {name: "empty", files: nil, want: map[string]string{}}, {name: "pairs", files: []string{"a:b", "c:d"}, want: map[string]string{"a": "b", "c": "d"}}, - {name: "drops unpaired", files: []string{"a", "c:d"}, want: map[string]string{"c": "d"}}, + {name: "rejects unpaired", files: []string{"a", "c:d"}, wantErr: true}, + {name: "rejects a second colon", files: []string{"a:b:c"}, wantErr: true}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got := SplitFiles(tt.files) + got, err := SplitFiles(tt.files) + if (err != nil) != tt.wantErr { + t.Fatalf("got %v, wantErr %v", err, tt.wantErr) + } + if tt.wantErr { + return + } if len(got) != len(tt.want) { t.Fatalf("got %v, want %v", got, tt.want) } diff --git a/cmd/utils/flags.go b/cmd/utils/flags.go index e0ff3f4..8762535 100644 --- a/cmd/utils/flags.go +++ b/cmd/utils/flags.go @@ -17,3 +17,11 @@ func FileFlag(usage string) *cli.StringSliceFlag { Usage: usage + ", src_path:dst_path[:mode[:uid:gid]]", } } + +func ForceFlag(usage string) *cli.BoolFlag { + return &cli.BoolFlag{ + Name: FlagForce, + Usage: usage, + Aliases: []string{"f"}, + } +} diff --git a/cmd/utils/resource.go b/cmd/utils/resource.go index bdd9ae9..ce71f51 100644 --- a/cmd/utils/resource.go +++ b/cmd/utils/resource.go @@ -14,6 +14,7 @@ const ( FlagExtraResources = "extra-resources" FlagFile = "file" + FlagForce = "force" ) // StorageParams builds the storage plugin request; zero values stay out, so an untouched entry defers to --extra-resources. diff --git a/cmd/utils/utils.go b/cmd/utils/utils.go index 2ace652..a05c9a3 100644 --- a/cmd/utils/utils.go +++ b/cmd/utils/utils.go @@ -11,7 +11,6 @@ import ( "strings" "text/template" - resourcetypes "github.com/projecteru2/core/resource/types" corepb "github.com/projecteru2/core/rpc/gen" "github.com/urfave/cli/v3" ) @@ -29,11 +28,6 @@ func GetNetworks(network string) map[string]string { return networks } -// ParseRAMInHuman parses a human-readable size ("100KB", "-1T") into bytes; the implementation lives with core's RawParams. -func ParseRAMInHuman(ram string) (int64, error) { - return resourcetypes.ParseRAMInHuman(ram) -} - // ParseDeployStrategy maps a --deploy-strategy value onto the core enum. func ParseDeployStrategy(name string) (corepb.DeployOptions_Strategy, error) { value, ok := corepb.DeployOptions_Strategy_value[strings.ToUpper(name)] diff --git a/cmd/utils/utils_test.go b/cmd/utils/utils_test.go index 7ee2d7f..b09bb6f 100644 --- a/cmd/utils/utils_test.go +++ b/cmd/utils/utils_test.go @@ -28,40 +28,6 @@ func TestGetNetworks(t *testing.T) { } } -func TestParseRAMInHuman(t *testing.T) { - tests := []struct { - name string - ram string - want int64 - wantErr bool - }{ - {name: "empty", ram: "", want: 0}, - {name: "bytes", ram: "1024", want: 1024}, - {name: "kilobytes", ram: "100KB", want: 102400}, - {name: "gigabytes", ram: "1G", want: 1 << 30}, - {name: "negative", ram: "-10G", want: -(10 << 30)}, - {name: "garbage", ram: "abc", wantErr: true}, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - got, err := ParseRAMInHuman(tt.ram) - if tt.wantErr { - if err == nil { - t.Fatalf("got %d, want an error", got) - } - return - } - if err != nil { - t.Fatalf("ParseRAMInHuman: %v", err) - } - if got != tt.want { - t.Errorf("got %d, want %d", got, tt.want) - } - }) - } -} - func TestSplitEquality(t *testing.T) { tests := []struct { name string diff --git a/cmd/workload/cmd.go b/cmd/workload/cmd.go index 9f73f49..caa3d8a 100644 --- a/cmd/workload/cmd.go +++ b/cmd/workload/cmd.go @@ -16,7 +16,6 @@ const ( flagEntry = "entry" flagNode = "node" flagPod = "pod" - flagForce = "force" flagEnv = "env" flagImage = "image" flagNetwork = "network" @@ -152,12 +151,7 @@ func Command() *cli.Command { ArgsUsage: workloadArgsUsage, Action: utils.ExitCoder(cmdWorkloadControl(corecluster.WorkloadStop)), Flags: []cli.Flag{ - &cli.BoolFlag{ - Name: flagForce, - Usage: "force to stop", - Aliases: []string{"f"}, - Value: false, - }, + utils.ForceFlag("force to stop"), }, }, { @@ -166,12 +160,7 @@ func Command() *cli.Command { ArgsUsage: workloadArgsUsage, Action: utils.ExitCoder(cmdWorkloadControl(corecluster.WorkloadStart)), Flags: []cli.Flag{ - &cli.BoolFlag{ - Name: flagForce, - Usage: "force to start", - Aliases: []string{"f"}, - Value: false, - }, + utils.ForceFlag("force to start"), }, }, { @@ -180,12 +169,7 @@ func Command() *cli.Command { ArgsUsage: workloadArgsUsage, Action: utils.ExitCoder(cmdWorkloadControl(corecluster.WorkloadRestart)), Flags: []cli.Flag{ - &cli.BoolFlag{ - Name: flagForce, - Usage: "force to restart", - Aliases: []string{"f"}, - Value: false, - }, + utils.ForceFlag("force to restart"), }, }, { @@ -194,12 +178,7 @@ func Command() *cli.Command { ArgsUsage: workloadArgsUsage, Action: utils.ExitCoder(cmdWorkloadRemove), Flags: []cli.Flag{ - &cli.BoolFlag{ - Name: flagForce, - Usage: "force to remove", - Aliases: []string{"f"}, - Value: false, - }, + utils.ForceFlag("force to remove"), }, }, { @@ -454,7 +433,6 @@ func Command() *cli.Command { &cli.Float64Flag{ Name: "cpu", Usage: "shortcut for cpu-request/limit, set them equally to this value", - Value: 1.0, }, &cli.StringFlag{ Name: flagMemoryRequest, @@ -469,7 +447,6 @@ func Command() *cli.Command { &cli.StringFlag{ Name: "memory", Usage: "shortcut for memory-request/limit, set them equally to this value", - Value: "512M", }, &cli.StringFlag{ Name: flagStorageRequest, diff --git a/cmd/workload/control.go b/cmd/workload/control.go index 4826831..b7c3d97 100644 --- a/cmd/workload/control.go +++ b/cmd/workload/control.go @@ -58,7 +58,7 @@ func cmdWorkloadControl(action string) cli.ActionFunc { client: client, ids: ids, action: action, - force: cmd.Bool(flagForce), + force: cmd.Bool(utils.FlagForce), } return o.run(ctx) } diff --git a/cmd/workload/copy.go b/cmd/workload/copy.go index 8faa7c4..0deca20 100644 --- a/cmd/workload/copy.go +++ b/cmd/workload/copy.go @@ -66,13 +66,20 @@ func (o *copyWorkloadsOptions) run(ctx context.Context) error { for filename, content := range files { storePath := filepath.Join(o.dir, filename) - if _, err := os.Stat(storePath); err == nil { - errs = errors.Join(errs, fmt.Errorf("%s already exists", storePath)) + f, err := os.OpenFile(storePath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600) //nolint:gosec + if err != nil { + if errors.Is(err, os.ErrExist) { + err = fmt.Errorf("%s already exists", storePath) + } + errs = errors.Join(errs, err) continue } - if err := os.WriteFile(storePath, content, 0o600); err != nil { + if _, err := f.Write(content); err != nil { errs = errors.Join(errs, fmt.Errorf("write %s: %w", storePath, err)) } + if err := f.Close(); err != nil { + errs = errors.Join(errs, fmt.Errorf("close %s: %w", storePath, err)) + } } return errs } diff --git a/cmd/workload/remove.go b/cmd/workload/remove.go index 1fdeaa1..465dfe5 100644 --- a/cmd/workload/remove.go +++ b/cmd/workload/remove.go @@ -51,7 +51,7 @@ func cmdWorkloadRemove(ctx context.Context, cmd *cli.Command) error { return err } - force := cmd.Bool(flagForce) + force := cmd.Bool(utils.FlagForce) if force { log.WithFunc("workload.cmdWorkloadRemove").Warn(ctx, "if workload not stopped, force to remove will not trigger hook process if set") } diff --git a/cmd/workload/replace.go b/cmd/workload/replace.go index 48f64ff..0cbe6f6 100644 --- a/cmd/workload/replace.go +++ b/cmd/workload/replace.go @@ -38,6 +38,11 @@ func cmdWorkloadReplace(ctx context.Context, cmd *cli.Command) error { return err } + copies, err := utils.SplitFiles(cmd.StringSlice("copy")) + if err != nil { + return err + } + networkInherit := cmd.Bool("network-inherit") if len(opts.Networks) > 0 { log.WithFunc("workload.cmdWorkloadReplace").Warn(ctx, "network is not empty, so network-inherit is set to false") @@ -46,7 +51,7 @@ func cmdWorkloadReplace(ctx context.Context, cmd *cli.Command) error { o := &replaceWorkloadsOptions{ client: client, opts: opts, - copies: utils.SplitFiles(cmd.StringSlice("copy")), + copies: copies, labels: utils.SplitEquality(cmd.StringSlice("label")), networkInherit: networkInherit, } diff --git a/cmd/workload/sendlarge.go b/cmd/workload/sendlarge.go index 9401882..d0ecbd2 100644 --- a/cmd/workload/sendlarge.go +++ b/cmd/workload/sendlarge.go @@ -32,7 +32,6 @@ func (o *sendLargeWorkloadsOptions) run(ctx context.Context) error { defer cancel() stream, err := o.client.SendLargeFile(ctx) if err != nil { - logger.Errorf(ctx, err, "send %s failed", o.dst) return err } @@ -68,7 +67,6 @@ func (o *sendLargeWorkloadsOptions) run(ctx context.Context) error { Owner: o.owners, Chunk: chunk[:n], }); err != nil { - logger.Errorf(ctx, err, "send %s failed", o.dst) wg.Wait() if errors.Is(err, io.EOF) { err = nil diff --git a/cmd/workload/utils.go b/cmd/workload/utils.go index 8f06af6..e6d64f6 100644 --- a/cmd/workload/utils.go +++ b/cmd/workload/utils.go @@ -128,18 +128,19 @@ func baseDeployOptions(ctx context.Context, cmd *cli.Command) (*corepb.DeployOpt } func ramOption(cmd *cli.Command, request, limit, shortcut string) (int64, int64, error) { - req, err := utils.ParseRAMInHuman(cmd.String(request)) + if cmd.IsSet(shortcut) { + both, err := resourcetypes.ParseRAMInHuman(cmd.String(shortcut)) + return both, both, err + } + + req, err := resourcetypes.ParseRAMInHuman(cmd.String(request)) if err != nil { return 0, 0, err } - lim, err := utils.ParseRAMInHuman(cmd.String(limit)) + lim, err := resourcetypes.ParseRAMInHuman(cmd.String(limit)) if err != nil { return 0, 0, err } - if cmd.IsSet(shortcut) { - both, err := utils.ParseRAMInHuman(cmd.String(shortcut)) - return both, both, err - } return req, lim, nil } diff --git a/describe/node.go b/describe/node.go index 7cec357..1e29d33 100644 --- a/describe/node.go +++ b/describe/node.go @@ -11,6 +11,8 @@ import ( corepb "github.com/projecteru2/core/rpc/gen" ) +type NodeResourceFilter func(cpumem, storage map[string]float64) bool + // Nodes describes nodes a command already holds. func Nodes(showInfo bool, nodes ...*corepb.Node) { describeOr(nodes, func(all []*corepb.Node) { renderNodes(showInfo, all...) }) @@ -18,16 +20,16 @@ func Nodes(showInfo bool, nodes ...*corepb.Node) { // NodesStream describes nodes as they arrive, one table per node when stream is set. func NodesStream(nodes <-chan *corepb.Node, showInfo, stream bool) { - describeChOr(nodes, func(ch <-chan *corepb.Node) { describeNodes(ch, showInfo, stream) }) + describeChOr(nodes, stream, func(all ...*corepb.Node) { renderNodes(showInfo, all...) }) } // NodeResource describes one node's resource. func NodeResource(ctx context.Context, resource *corepb.NodeResource) { - describeOr(resource, func(r *corepb.NodeResource) { renderNodeResources(ctx, r) }) + describeOr(resource, func(r *corepb.NodeResource) { renderNodeResources(nodePercents(ctx, r)...) }) } -func NodeResources(ctx context.Context, resources <-chan *corepb.NodeResource, stream bool) { - describeChOr(resources, func(ch <-chan *corepb.NodeResource) { describeNodeResources(ctx, ch, stream) }) +func NodeResources(ctx context.Context, resources <-chan *corepb.NodeResource, stream bool, keep NodeResourceFilter) { + describeChOr(nodePercentChan(ctx, resources, keep), stream, renderNodeResources) } // NodeStatusMessage describes node status messages as json, yaml or log lines. @@ -35,21 +37,6 @@ func NodeStatusMessage(ctx context.Context, ms ...*corepb.NodeStatusStreamMessag describeOr(ms, func(m []*corepb.NodeStatusStreamMessage) { describeNodeStatusMessage(ctx, m) }) } -func describeNodes(nodes <-chan *corepb.Node, showInfo, stream bool) { - if stream { - for node := range nodes { - renderNodes(showInfo, node) - } - return - } - - all := []*corepb.Node{} - for node := range nodes { - all = append(all, node) - } - renderNodes(showInfo, all...) -} - func renderNodes(showInfo bool, nodes ...*corepb.Node) { capacities := make([]resourcetypes.Resources, len(nodes)) usages := make([]resourcetypes.Resources, len(nodes)) @@ -95,35 +82,50 @@ func nodePluginRows(capacity, usage resourcetypes.RawParams) []string { return append(rows, parseAll(usage)...) } -func describeNodeResources(ctx context.Context, resources <-chan *corepb.NodeResource, stream bool) { - if stream { - for resource := range resources { - renderNodeResources(ctx, resource) - } - return - } - all := []*corepb.NodeResource{} - for resource := range resources { - all = append(all, resource) - } - renderNodeResources(ctx, all...) +type nodePercent struct { + *corepb.NodeResource + cpumem map[string]float64 + storage map[string]float64 } -func renderNodeResources(ctx context.Context, resources ...*corepb.NodeResource) { - logger := log.WithFunc("describe.renderNodeResources") - groups := make([][][]string, 0, len(resources)) +func nodePercents(ctx context.Context, resources ...*corepb.NodeResource) []nodePercent { + logger := log.WithFunc("describe.nodePercents") + rv := make([]nodePercent, 0, len(resources)) for _, resource := range resources { cr, sr, err := ToResourcePercent(resource) if err != nil { logger.Errorf(ctx, err, "resource percent of node %s", resource.Name) continue } + rv = append(rv, nodePercent{resource, cr, sr}) + } + return rv +} + +func nodePercentChan(ctx context.Context, resources <-chan *corepb.NodeResource, keep NodeResourceFilter) <-chan nodePercent { + rv := make(chan nodePercent) + go func() { + defer close(rv) + for resource := range resources { + for _, percent := range nodePercents(ctx, resource) { + if keep == nil || keep(percent.cpumem, percent.storage) { + rv <- percent + } + } + } + }() + return rv +} + +func renderNodeResources(resources ...nodePercent) { + groups := make([][][]string, 0, len(resources)) + for _, resource := range resources { groups = append(groups, [][]string{ {resource.Name}, - {fmt.Sprintf("%.2f%%", cr["cpu"]*100)}, - {fmt.Sprintf("%.2f%%", cr["memory"]*100)}, - {fmt.Sprintf("%.2f%%", sr["storage"]*100)}, - {fmt.Sprintf("%.2f%%", sr["volumes"]*100)}, + {fmt.Sprintf("%.2f%%", resource.cpumem["cpu"]*100)}, + {fmt.Sprintf("%.2f%%", resource.cpumem["memory"]*100)}, + {fmt.Sprintf("%.2f%%", resource.storage["storage"]*100)}, + {fmt.Sprintf("%.2f%%", resource.storage["volumes"]*100)}, {strings.Join(resource.Diffs, "\n")}, }) } diff --git a/describe/node_test.go b/describe/node_test.go index 44da567..a2b006a 100644 --- a/describe/node_test.go +++ b/describe/node_test.go @@ -157,7 +157,7 @@ func TestNodeResources(t *testing.T) { Format = tt.format t.Cleanup(func() { Format = "" }) - got := captureStdout(t, func() { NodeResources(t.Context(), ToChan(testNodeResources()...), false) }) + got := captureStdout(t, func() { NodeResources(t.Context(), ToChan(testNodeResources()...), false, nil) }) if got != tt.want { t.Errorf("got\n%s\nwant\n%s", got, tt.want) } @@ -165,6 +165,26 @@ func TestNodeResources(t *testing.T) { } } +func TestNodeResourcesKeepsOnlyMatchingNodes(t *testing.T) { + resources := []*corepb.NodeResource{ + {Name: "unparsable", ResourceUsage: "{"}, + {Name: "idle", ResourceUsage: `{"cpumem":{"cpu":1}}`, ResourceCapacity: `{"cpumem":{"cpu":8}}`}, + {Name: "busy", ResourceUsage: `{"cpumem":{"cpu":6}}`, ResourceCapacity: `{"cpumem":{"cpu":8}}`}, + } + keep := func(cpumem, _ map[string]float64) bool { return cpumem["cpu"] > 0.5 } + + want := `┌──────┬────────┬────────┬─────────┬────────┬───────┐ +│ NAME │ CPU │ MEMORY │ STORAGE │ VOLUME │ DIFFS │ +├──────┼────────┼────────┼─────────┼────────┼───────┤ +│ busy │ 75.00% │ 0.00% │ 0.00% │ 0.00% │ │ +└──────┴────────┴────────┴─────────┴────────┴───────┘ +` + got := captureStdout(t, func() { NodeResources(t.Context(), ToChan(resources...), false, keep) }) + if got != want { + t.Errorf("got\n%s\nwant\n%s", got, want) + } +} + func TestNodesStream(t *testing.T) { want := `┌───────┬─────────────────────┬────────────────┬──────────────┐ │ NAME │ ENDPOINT │ STATUS │ CPUMEM │ diff --git a/describe/utils.go b/describe/utils.go index b581057..e2d8f93 100644 --- a/describe/utils.go +++ b/describe/utils.go @@ -37,44 +37,31 @@ func ToResourcePercent(resource *corepb.NodeResource) (cpumem, storage map[strin storageCap := resCap[utils.ResourceStorage] cr, sr := map[string]float64{}, map[string]float64{} if cpumemUsage != nil && cpumemCap != nil { - cpuUsage := cpumemUsage.Float64("cpu") - cpuCap := cpumemCap.Float64("cpu") - memUsage := cpumemUsage.Float64("memory") - memCap := cpumemCap.Float64("memory") - cr["cpu"] = 0.0 - cr["memory"] = 0.0 - if cpuCap != 0 { - cr["cpu"] = cpuUsage / cpuCap - } - if memCap != 0 { - cr["memory"] = memUsage / memCap - } + cr["cpu"] = ratio(cpumemUsage.Float64("cpu"), cpumemCap.Float64("cpu")) + cr["memory"] = ratio(cpumemUsage.Float64("memory"), cpumemCap.Float64("memory")) } if storageUsage != nil && storageCap != nil { - stUsage := storageUsage.Float64("storage") - stCap := storageCap.Float64("storage") - volumesUsage := storageUsage.RawParams("volumes") - volumesCap := storageCap.RawParams("volumes") - sr["storage"] = 0.0 - sr["volumes"] = 0.0 - if stCap != 0 { - sr["storage"] = stUsage / stCap - } - vu := 0.0 - vc := 0.0 - for k := range volumesUsage { - vu += volumesUsage.Float64(k) - } - for k := range volumesCap { - vc += volumesCap.Float64(k) - } - if vc != 0 { - sr["volumes"] = vu / vc - } + sr["storage"] = ratio(storageUsage.Float64("storage"), storageCap.Float64("storage")) + sr["volumes"] = ratio(sumParams(storageUsage.RawParams("volumes")), sumParams(storageCap.RawParams("volumes"))) } return cr, sr, nil } +func ratio(usage, capacity float64) float64 { + if capacity == 0 { + return 0 + } + return usage / capacity +} + +func sumParams(params resourcetypes.RawParams) float64 { + sum := 0.0 + for key := range params { + sum += params.Float64(key) + } + return sum +} + func isJSON() bool { return strings.ToLower(Format) == "json" } @@ -158,9 +145,7 @@ func pluginNames(resourceSets ...[]resourcetypes.Resources) []string { func unmarshalResources(encoded string) resourcetypes.Resources { res := resourcetypes.Resources{} - if len(encoded) > 0 { - _ = json.Unmarshal([]byte(encoded), &res) - } + _ = json.Unmarshal([]byte(encoded), &res) return res } @@ -185,7 +170,7 @@ func describeOr[T any](v T, fallback func(T)) { } } -func describeChOr[T any](ch <-chan T, fallback func(<-chan T)) { +func describeChOr[T any](ch <-chan T, stream bool, render func(...T)) { collect := func() []T { items := []T{} for t := range ch { @@ -198,8 +183,12 @@ func describeChOr[T any](ch <-chan T, fallback func(<-chan T)) { describeAsJSON(collect()) case isYAML(): describeAsYAML(collect()) + case stream: + for t := range ch { + render(t) + } default: - fallback(ch) + render(collect()...) } } diff --git a/describe/workload.go b/describe/workload.go index acc4100..f03275b 100644 --- a/describe/workload.go +++ b/describe/workload.go @@ -48,7 +48,7 @@ func WorkloadStatuses(workloadStatuses ...*corepb.WorkloadStatus) { func describeStatistics(stat workloadStatistics) { renderTable([]string{"CPUs", "Memory", "Storage"}, [][]string{ - {fmt.Sprintf("%f", stat.CPUs)}, + {strconv.FormatFloat(stat.CPUs, 'f', -1, 64)}, {strconv.FormatInt(stat.Memory, 10)}, {strconv.FormatInt(stat.Storage, 10)}, }) @@ -84,9 +84,7 @@ func workloadNetworks(workload *corepb.Workload) []string { maps.Copy(addresses, workload.Status.Networks) } for name := range workload.Publish { - if _, ok := addresses[name]; !ok { - addresses[name] = "" - } + addresses[name] = "" } ns := []string{} diff --git a/describe/workload_test.go b/describe/workload_test.go index 3b4e173..388b8e2 100644 --- a/describe/workload_test.go +++ b/describe/workload_test.go @@ -147,11 +147,11 @@ storage: 2048 workloads: []*corepb.Workload{ {Resources: `{"cpumem":{"cpu_request":1.5,"memory_request":1024},"resource-storage":{"storage_request":2048}}`}, }, - want: `┌──────────┬────────┬─────────┐ -│ CPUS │ MEMORY │ STORAGE │ -├──────────┼────────┼─────────┤ -│ 1.500000 │ 1024 │ 2048 │ -└──────────┴────────┴─────────┘ + want: `┌──────┬────────┬─────────┐ +│ CPUS │ MEMORY │ STORAGE │ +├──────┼────────┼─────────┤ +│ 1.5 │ 1024 │ 2048 │ +└──────┴────────┴─────────┘ `, }, } diff --git a/interactive/stream.go b/interactive/stream.go index 4962155..85a8696 100644 --- a/interactive/stream.go +++ b/interactive/stream.go @@ -131,12 +131,11 @@ func attachTerminal(ctx context.Context, iStream Stream) func() { ctx, cancel := context.WithCancel(ctx) state, err := term.MakeRaw(stdinFd) + go pumpStdin(ctx, iStream) if err != nil { // stdin is a pipe or a file: no raw mode and no window size to report. - go pumpStdin(ctx, iStream) return cancel } - go pumpStdin(ctx, iStream) sigs := make(chan os.Signal, 1) signal.Notify(sigs, syscall.SIGWINCH) diff --git a/types/specs.go b/types/specs.go index a70da8c..6e5583c 100644 --- a/types/specs.go +++ b/types/specs.go @@ -8,13 +8,13 @@ import ( // Specs is the deploy spec file of an application. type Specs struct { - Appname string `yaml:"appname,omitempty"` - Entrypoints map[string]Entrypoint `yaml:"entrypoints,omitempty,flow"` - Volumes []string `yaml:"volumes,omitempty,flow"` - VolumesRequest []string `yaml:"volumes_request,omitempty,flow"` - Labels map[string]string `yaml:"labels,omitempty,flow"` - DNS []string `yaml:"dns,omitempty,flow"` - ExtraHosts []string `yaml:"extra_hosts,omitempty,flow"` + Appname string `yaml:"appname"` + Entrypoints map[string]Entrypoint `yaml:"entrypoints"` + Volumes []string `yaml:"volumes"` + VolumesRequest []string `yaml:"volumes_request"` + Labels map[string]string `yaml:"labels"` + DNS []string `yaml:"dns"` + ExtraHosts []string `yaml:"extra_hosts"` } // Entrypoint accepts both the legacy `cmd` string and the current `commands` list.