-
Notifications
You must be signed in to change notification settings - Fork 3.1k
/
watch.go
117 lines (105 loc) · 2.48 KB
/
watch.go
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
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
package commands
import (
"fmt"
"os"
"time"
"github.com/argoproj/pkg/errors"
"github.com/spf13/cobra"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/fields"
"github.com/argoproj/argo/cmd/argo/commands/client"
workflowpkg "github.com/argoproj/argo/pkg/apiclient/workflow"
wfv1 "github.com/argoproj/argo/pkg/apis/workflow/v1alpha1"
"github.com/argoproj/argo/workflow/packer"
)
func NewWatchCommand() *cobra.Command {
var command = &cobra.Command{
Use: "watch WORKFLOW",
Short: "watch a workflow until it completes",
Run: func(cmd *cobra.Command, args []string) {
if len(args) != 1 {
cmd.HelpFunc()(cmd, args)
os.Exit(1)
}
watchWorkflow(args[0])
},
}
return command
}
func apiServerWatchWorkflow(wfName string) {
conn := client.GetClientConn()
defer conn.Close()
apiClient, ctx := GetWFApiServerGRPCClient(conn)
fieldSelector := fields.ParseSelectorOrDie(fmt.Sprintf("metadata.name=%s", wfName))
wfReq := workflowpkg.WatchWorkflowsRequest{
Namespace: namespace,
ListOptions: &metav1.ListOptions{
FieldSelector: fieldSelector.String(),
},
}
stream, err := apiClient.WatchWorkflows(ctx, &wfReq)
if err != nil {
errors.CheckError(err)
return
}
for {
event, err := stream.Recv()
if err != nil {
errors.CheckError(err)
break
}
wf := event.Object
if wf != nil {
printWorkflowStatus(wf)
if !wf.Status.FinishedAt.IsZero() {
break
}
} else {
break
}
}
}
func watchWorkflow(name string) {
if client.ArgoServer != "" {
apiServerWatchWorkflow(name)
} else {
InitWorkflowClient()
k8sApiWatchWorkflow(name)
}
}
func k8sApiWatchWorkflow(name string) {
fieldSelector := fields.ParseSelectorOrDie(fmt.Sprintf("metadata.name=%s", name))
opts := metav1.ListOptions{
FieldSelector: fieldSelector.String(),
}
wf, err := wfClient.Get(name, metav1.GetOptions{})
errors.CheckError(err)
watchIf, err := wfClient.Watch(opts)
errors.CheckError(err)
ticker := time.NewTicker(time.Second)
for {
select {
case next := <-watchIf.ResultChan():
wf, _ = next.Object.(*wfv1.Workflow)
case <-ticker.C:
}
if wf == nil {
watchIf.Stop()
watchIf, err = wfClient.Watch(opts)
errors.CheckError(err)
continue
}
printWorkflowStatus(wf)
if !wf.Status.FinishedAt.IsZero() {
break
}
}
watchIf.Stop()
}
func printWorkflowStatus(wf *wfv1.Workflow) {
err := packer.DecompressWorkflow(wf)
errors.CheckError(err)
print("\033[H\033[2J")
print("\033[0;0H")
printWorkflowHelper(wf, getFlags{})
}