Repository navigation
Expand file tree
/
Copy pathtask_queue_commands.go
More file actions
99 lines (87 loc) · 2.73 KB
/
Copy pathtask_queue_commands.go
File metadata and controls
99 lines (87 loc) · 2.73 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
package taskqueue
import (
"fmt"
"strings"
"github.com/temporalio/cli/client"
"github.com/temporalio/cli/common"
"github.com/temporalio/tctl-kit/pkg/color"
"github.com/temporalio/tctl-kit/pkg/output"
"github.com/urfave/cli/v2"
enumspb "go.temporal.io/api/enums/v1"
taskqueuepb "go.temporal.io/api/taskqueue/v1"
"go.temporal.io/api/workflowservice/v1"
)
// DescribeTaskQueue show pollers info of a given taskqueue
func DescribeTaskQueue(c *cli.Context) error {
sdkClient, err := client.GetSDKClient(c)
if err != nil {
return err
}
taskQueue := c.String(common.FlagTaskQueue)
taskQueueType := strToTaskQueueType(c.String(common.FlagTaskQueueType))
ctx, cancel := common.NewContext(c)
defer cancel()
resp, err := sdkClient.DescribeTaskQueue(ctx, taskQueue, taskQueueType)
if err != nil {
return fmt.Errorf("unable to describe task queue: %w", err)
}
opts := &output.PrintOptions{
// TODO enable when versioning feature is out
// Fields: []string{"Identity", "LastAccessTime", "RatePerSecond", "WorkerVersioningId"},
Fields: []string{"Identity", "LastAccessTime", "RatePerSecond"},
}
var items []interface{}
for _, e := range resp.Pollers {
items = append(items, e)
}
return output.PrintItems(c, items, opts)
}
// ListTaskQueuePartitions gets all the taskqueue partition and host information.
func ListTaskQueuePartitions(c *cli.Context) error {
frontendClient := client.CFactory.FrontendClient(c)
namespace, err := common.RequiredFlag(c, common.FlagNamespace)
if err != nil {
return err
}
taskQueue := c.String(common.FlagTaskQueue)
ctx, cancel := common.NewContext(c)
defer cancel()
request := &workflowservice.ListTaskQueuePartitionsRequest{
Namespace: namespace,
TaskQueue: &taskqueuepb.TaskQueue{
Name: taskQueue,
Kind: enumspb.TASK_QUEUE_KIND_NORMAL,
},
}
resp, err := frontendClient.ListTaskQueuePartitions(ctx, request)
if err != nil {
return fmt.Errorf("unable to list task queues: %w", err)
}
optsW := &output.PrintOptions{
Fields: []string{"Key", "OwnerHostName"},
}
var items []interface{}
fmt.Println(color.Magenta(c, "Workflow Task Queue Partitions\n"))
for _, e := range resp.WorkflowTaskQueuePartitions {
items = append(items, e)
}
err = output.PrintItems(c, items, optsW)
if err != nil {
return err
}
optsA := &output.PrintOptions{
Fields: []string{"Key", "OwnerHostName"},
}
items = items[:0]
fmt.Println(color.Magenta(c, "\nActivity Task Queue Partitions\n"))
for _, e := range resp.ActivityTaskQueuePartitions {
items = append(items, e)
}
return output.PrintItems(c, items, optsA)
}
func strToTaskQueueType(str string) enumspb.TaskQueueType {
if strings.ToLower(str) == "activity" {
return enumspb.TASK_QUEUE_TYPE_ACTIVITY
}
return enumspb.TASK_QUEUE_TYPE_WORKFLOW
}