AdminListTaskQueueTasks displays task information
(c *cli.Context, clientFactory ClientFactory)
| 13 | |
| 14 | // AdminListTaskQueueTasks displays task information |
| 15 | func AdminListTaskQueueTasks(c *cli.Context, clientFactory ClientFactory) error { |
| 16 | namespace, err := getRequiredOption(c, FlagNamespace) |
| 17 | if err != nil { |
| 18 | return err |
| 19 | } |
| 20 | tqName := c.String(FlagTaskQueue) |
| 21 | tlTypeInt, err := StringToEnum(c.String(FlagTaskQueueType), enumspb.TaskQueueType_value) |
| 22 | if err != nil { |
| 23 | return fmt.Errorf("invalid task queue type: %v", err) |
| 24 | } |
| 25 | tqType := enumspb.TaskQueueType(tlTypeInt) |
| 26 | if tqType == enumspb.TASK_QUEUE_TYPE_UNSPECIFIED { |
| 27 | return fmt.Errorf("missing Task Queue type") |
| 28 | } |
| 29 | minTaskID := c.Int64(FlagMinTaskID) |
| 30 | maxTaskID := c.Int64(FlagMaxTaskID) |
| 31 | pageSize := c.Int(FlagPageSize) |
| 32 | workflowID := c.String(FlagWorkflowID) |
| 33 | runID := c.String(FlagRunID) |
| 34 | subqueue := c.Int(FlagSubqueue) |
| 35 | var minPass int64 |
| 36 | if c.Bool(FlagFair) { |
| 37 | minPass = c.Int64(FlagMinPass) |
| 38 | } else if c.IsSet(FlagMinPass) { |
| 39 | return fmt.Errorf("flag --%s is only valid with --%s", FlagMinPass, FlagFair) |
| 40 | } |
| 41 | client := clientFactory.AdminClient(c) |
| 42 | |
| 43 | req := &adminservice.GetTaskQueueTasksRequest{ |
| 44 | Namespace: namespace, |
| 45 | TaskQueue: tqName, |
| 46 | TaskQueueType: tqType, |
| 47 | MinTaskId: minTaskID, |
| 48 | MaxTaskId: maxTaskID, |
| 49 | BatchSize: int32(pageSize), |
| 50 | Subqueue: int32(subqueue), |
| 51 | MinPass: minPass, |
| 52 | } |
| 53 | |
| 54 | paginationFunc := func(paginationToken []byte) ([]any, []byte, error) { |
| 55 | ctx, cancel := newContext(c) |
| 56 | defer cancel() |
| 57 | |
| 58 | req.NextPageToken = paginationToken |
| 59 | response, err := client.GetTaskQueueTasks(ctx, req) |
| 60 | if err != nil { |
| 61 | return nil, nil, err |
| 62 | } |
| 63 | |
| 64 | tasks := response.Tasks |
| 65 | if workflowID != "" { |
| 66 | filteredTasks := tasks[:0] |
| 67 | |
| 68 | for _, task := range tasks { |
| 69 | if task.Data.WorkflowId != workflowID { |
| 70 | continue |
| 71 | } |
| 72 | if runID != "" && task.Data.RunId != runID { |
no test coverage detected