/
job.go
64 lines (47 loc) · 1.3 KB
/
job.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
package k8s
import (
"context"
"errors"
"fmt"
"io"
log "github.com/sirupsen/logrus"
corev1 "k8s.io/api/core/v1"
"k8s.io/client-go/kubernetes"
applog "github.com/utkuozdemir/pv-migrate/internal/log"
"github.com/utkuozdemir/pv-migrate/internal/rsync"
)
var ErrJobFailed = errors.New("job failed")
func WaitForJobCompletion(logger *log.Entry, cli kubernetes.Interface,
namespace string, name string, progressBarRequested bool,
) error {
s := fmt.Sprintf("job-name=%s", name)
pod, err := WaitForPod(cli, namespace, s)
if err != nil {
return err
}
successCh := make(chan bool, 1)
showProgressBar := progressBarRequested &&
logger.Context.Value(applog.FormatContextKey) == applog.FormatFancy
logTail := rsync.LogTail{
LogReaderFunc: func() (io.ReadCloser, error) {
return cli.CoreV1().Pods(namespace).GetLogs(pod.Name,
&corev1.PodLogOptions{Follow: true}).Stream(context.TODO())
},
SuccessCh: successCh,
ShowProgressBar: showProgressBar,
Logger: logger,
}
go logTail.Start()
terminatedPod, err := waitForPodTermination(cli, pod.Namespace, pod.Name)
if err != nil {
successCh <- false
return err
}
if *terminatedPod != corev1.PodSucceeded {
successCh <- false
err := fmt.Errorf("%w: %s", ErrJobFailed, name)
return err
}
successCh <- true
return nil
}