This repository has been archived by the owner on Jun 5, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 1
/
operation.go
75 lines (57 loc) · 1.48 KB
/
operation.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
package mapreduce
import (
"context"
"regexp"
"time"
"github.com/go-faster/yt/yt"
"github.com/go-faster/yt/yterrors"
)
type Operation interface {
ID() yt.OperationID
Wait() error
}
type operation struct {
yc yt.Client
ctx context.Context
opID yt.OperationID
}
func (o *operation) ID() yt.OperationID {
return o.opID
}
var failedJobLimitExceededRE = regexp.MustCompile("Failed jobs limit exceeded")
func (o *operation) getOperationError(status *yt.OperationStatus) error {
innerErr := status.Result.Error
if yterrors.ContainsMessageRE(status.Result.Error, failedJobLimitExceededRE) {
result, err := o.yc.ListJobs(o.ctx, o.opID, &yt.ListJobsOptions{JobState: &yt.JobFailed})
if err != nil {
return yterrors.Err("unable to get list of failed jobs", innerErr)
}
if len(result.Jobs) == 0 {
return yterrors.Err("no failed jobs found", innerErr)
}
job := result.Jobs[0]
stderr, err := o.yc.GetJobStderr(o.ctx, o.opID, job.ID, nil)
if err != nil {
return yterrors.Err("unable to get job stderr", innerErr)
}
return yterrors.Err("job failed",
innerErr,
yterrors.Attr("stderr", string(stderr)))
}
return status.Result.Error
}
func (o *operation) Wait() error {
for {
status, err := o.yc.GetOperation(o.ctx, o.opID, nil)
if err != nil {
return err
}
if status.State.IsFinished() {
if status.Result.Error != nil && status.Result.Error.Code != 0 {
return o.getOperationError(status)
}
return nil
}
time.Sleep(time.Second * 5)
}
}