• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

opendefensecloud / artifact-conduit / 37336100766

05 Oct 2026 03:50PM UTC coverage: 85.185%. First build
37336100766

Pull #475

github

web-flow
Merge 18d612576 into 0564c7166
Pull Request #475: chore(deps): update golang version sync

9 of 11 new or added lines in 3 files covered. (81.82%)

1357 of 1593 relevant lines covered (85.19%)

1817.89 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

84.55
/pkg/controller/workflow_handler.go
1
// Copyright BWI GmbH and Artifact Conduit contributors
2
// SPDX-License-Identifier: Apache-2.0
3

4
package controller
5

6
import (
7
        "context"
8
        "fmt"
9

10
        wfv1alpha1 "github.com/argoproj/argo-workflows/v4/pkg/apis/workflow/v1alpha1"
11
        "github.com/go-logr/logr"
12
        corev1 "k8s.io/api/core/v1"
13
        "sigs.k8s.io/controller-runtime/pkg/client"
14
        "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
15

16
        arcv1alpha1 "go.opendefense.cloud/arc/api/arc/v1alpha1"
17
        "go.opendefense.cloud/arc/pkg/metrics"
18
)
19

20
type WorkflowHandler interface {
21
        DeleteArgoResources(ctx context.Context) error
22
        CreateArgoResources(ctx context.Context) error
23
        CheckArgoResources(ctx context.Context) error
24
}
25

26
var _ WorkflowHandler = &SingleWorkflowHandler{}
27

28
type SingleWorkflowHandler struct {
29
        *ArtifactWorkflowReconciler
30
        log logr.Logger
31
        aw  *arcv1alpha1.ArtifactWorkflow
32
}
33

34
func NewSingleWorkflowHandler(r *ArtifactWorkflowReconciler, log logr.Logger, aw *arcv1alpha1.ArtifactWorkflow) *SingleWorkflowHandler {
35
        return &SingleWorkflowHandler{r, log, aw}
4,039 ✔
36
}
4,039 ✔
37

38
func (h *SingleWorkflowHandler) DeleteArgoResources(ctx context.Context) error {
39
        wf := wfv1alpha1.Workflow{
9 ✔
40
                Namespace: h.aw.Namespace,
9 ✔
41
                Name:      h.aw.Name,
9 ✔
42
        }
9 ✔
43
        if err := h.Delete(ctx, &wf); client.IgnoreNotFound(err) != nil {
9 ✔
44
                h.Recorder.Eventf(h.aw, nil, corev1.EventTypeWarning, ReasonDeletionFailed, "Delete", fmt.Sprintf("Failed to delete associated workflow '%s': %v", h.aw.Name, err))
×
45
                metrics.RecordReconcileError(ControllerArtifactWorkflow, ReasonDeletionFailed)
×
46

×
47
                return errLogAndWrap(h.log, err, "workflow deletion failed")
×
48
        }
×
49
        h.Recorder.Eventf(h.aw, nil, corev1.EventTypeNormal, "Deleted", "Delete", fmt.Sprintf("Deleted workflow '%s'", h.aw.Name))
9 ✔
50

9 ✔
51
        return nil
9 ✔
52
}
53

54
func (h *SingleWorkflowHandler) CreateArgoResources(ctx context.Context) error {
55
        srcSecret, dstSecret, err := h.retrieveSecrets(ctx, h.aw)
2,131 ✔
56
        if err != nil {
2,131 ✔
57
                return errLogAndWrap(h.log, err, "failed to fetch secrets for artifact workflow")
13 ✔
58
        }
13 ✔
59

60
        wf := hydrateArgoWorkflow(h.aw, srcSecret, dstSecret)
2,118 ✔
61

2,118 ✔
62
        if err := controllerutil.SetControllerReference(h.aw, wf, h.Scheme); err != nil {
2,118 ✔
63
                return errLogAndWrap(h.log, err, "failed to set controller reference")
×
64
        }
×
65

66
        if err := h.Create(ctx, wf); client.IgnoreAlreadyExists(err) != nil {
2,118 ✔
67
                h.Recorder.Eventf(h.aw, nil, corev1.EventTypeWarning, ReasonCreationFailed, "Create", fmt.Sprintf("Failed to create workflow '%s': %v", wf.GetName(), err))
260 ✔
68
                metrics.RecordReconcileError(ControllerArtifactWorkflow, ReasonCreationFailed)
260 ✔
69

260 ✔
70
                return errLogAndWrap(h.log, err, "failed to create argo workflow")
260 ✔
71
        }
260 ✔
72
        h.Recorder.Eventf(h.aw, nil, corev1.EventTypeNormal, "Created", "Create", fmt.Sprintf("Created workflow '%s'", wf.GetName()))
1,858 ✔
73

1,858 ✔
74
        h.aw.Status.Phase = arcv1alpha1.WorkflowPending
1,858 ✔
75
        if err := h.Status().Update(ctx, h.aw); err != nil {
1,858 ✔
76
                return errLogAndWrap(h.log, err, "failed to update status")
3 ✔
77
        }
3 ✔
78

79
        return nil
1,855 ✔
80
}
81

82
func (h *SingleWorkflowHandler) CheckArgoResources(ctx context.Context) error {
83
        wf := wfv1alpha1.Workflow{}
1,858 ✔
84
        if err := h.Get(ctx, namespacedName(h.aw.Namespace, h.aw.Name), &wf); err != nil {
1,858 ✔
85
                return errLogAndWrap(h.log, err, "failed to get workflow")
×
86
        }
×
87

88
        updated, done := h.setStatusFromWorkflow(ctx, h.log, h.aw, &wf)
1,858 ✔
89

1,858 ✔
90
        if wf.Status.Phase == wfv1alpha1.WorkflowSucceeded && h.aw.Status.Succeeded != 1 {
1,858 ✔
91
                h.aw.Status.Succeeded = 1
4 ✔
92
                updated = true
4 ✔
93
        }
4 ✔
94

95
        failed := wf.Status.Phase == wfv1alpha1.WorkflowError ||
1,858 ✔
96
                wf.Status.Phase == wfv1alpha1.WorkflowFailed
1,858 ✔
97

1,858 ✔
98
        if failed && h.aw.Status.Failed != 1 {
1,858 ✔
99
                h.aw.Status.Failed = 1
3 ✔
100
                updated = true
3 ✔
101
        }
3 ✔
102

103
        if !updated {
1,858 ✔
104
                return nil
3 ✔
105
        }
3 ✔
106

107
        if err := h.Status().Update(ctx, h.aw); err != nil {
1,855 ✔
108
                return errLogAndWrap(h.log, err, "failed to update status")
3 ✔
109
        }
3 ✔
110

111
        if done != nil {
1,852 ✔
112
                done.record()
7 ✔
113
        }
7 ✔
114

115
        return nil
1,852 ✔
116
}
117

118
var _ WorkflowHandler = &CronWorkflowHandler{}
119

120
type CronWorkflowHandler struct {
121
        *ArtifactWorkflowReconciler
122
        log logr.Logger
123
        aw  *arcv1alpha1.ArtifactWorkflow
124
}
125

126
func NewCronWorkflowHandler(r *ArtifactWorkflowReconciler, log logr.Logger, aw *arcv1alpha1.ArtifactWorkflow) *CronWorkflowHandler {
127
        return &CronWorkflowHandler{r, log, aw}
31 ✔
128
}
31 ✔
129

130
func (h *CronWorkflowHandler) DeleteArgoResources(ctx context.Context) error {
131
        cwf := wfv1alpha1.CronWorkflow{
×
NEW
132
                Namespace: h.aw.Namespace,
×
NEW
133
                Name:      h.aw.Name,
×
134
        }
×
135
        if err := h.Delete(ctx, &cwf); client.IgnoreNotFound(err) != nil {
×
136
                h.Recorder.Eventf(h.aw, nil, corev1.EventTypeWarning, ReasonDeletionFailed, "Delete", fmt.Sprintf("Failed to delete associated cron workflow '%s': %v", h.aw.Name, err))
×
137
                metrics.RecordReconcileError(ControllerArtifactWorkflow, ReasonDeletionFailed)
×
138

×
139
                return errLogAndWrap(h.log, err, "cron workflow deletion failed")
×
140
        }
×
141
        h.Recorder.Eventf(h.aw, nil, corev1.EventTypeNormal, "Deleted", "Delete", fmt.Sprintf("Deleted cron workflow '%s'", h.aw.Name))
×
142

×
143
        return nil
×
144
}
145

146
func (h *CronWorkflowHandler) CreateArgoResources(ctx context.Context) error {
147
        srcSecret, dstSecret, err := h.retrieveSecrets(ctx, h.aw)
14 ✔
148
        if err != nil {
14 ✔
149
                return errLogAndWrap(h.log, err, "failed to fetch secrets for artifact workflow")
×
150
        }
×
151

152
        cwf := hydrateArgoCronWorkflow(h.aw, srcSecret, dstSecret)
14 ✔
153

14 ✔
154
        if err := controllerutil.SetControllerReference(h.aw, cwf, h.Scheme); err != nil {
14 ✔
155
                return errLogAndWrap(h.log, err, "failed to set controller reference")
×
156
        }
×
157

158
        if err := h.Create(ctx, cwf); err != nil {
14 ✔
159
                if client.IgnoreAlreadyExists(err) != nil {
13 ✔
160
                        h.Recorder.Eventf(h.aw, nil, corev1.EventTypeWarning, ReasonCreationFailed, "Create", fmt.Sprintf("Failed to create cron workflow '%s': %v", cwf.GetName(), err))
13 ✔
161
                        metrics.RecordReconcileError(ControllerArtifactWorkflow, ReasonCreationFailed)
13 ✔
162

13 ✔
163
                        return errLogAndWrap(h.log, err, "failed to create argo cron workflow")
13 ✔
164
                }
13 ✔
165
        } else {
166
                h.Recorder.Eventf(h.aw, nil, corev1.EventTypeNormal, "Created", "Create", fmt.Sprintf("Created cron workflow '%s'", cwf.GetName()))
1 ✔
167
        }
1 ✔
168

169
        h.aw.Status.Phase = arcv1alpha1.WorkflowPending
1 ✔
170
        if err := h.Status().Update(ctx, h.aw); err != nil {
1 ✔
171
                return errLogAndWrap(h.log, err, "failed to update status")
×
172
        }
×
173

174
        return nil
1 ✔
175
}
176

177
func (h *CronWorkflowHandler) CheckArgoResources(ctx context.Context) error {
178
        cwf := wfv1alpha1.CronWorkflow{}
15 ✔
179
        if err := h.Get(ctx, namespacedName(h.aw.Namespace, h.aw.Name), &cwf); err != nil {
15 ✔
180
                return errLogAndWrap(h.log, err, "failed to get cron workflow")
×
181
        }
×
182

183
        updated := false
15 ✔
184

15 ✔
185
        if !h.aw.Status.LastScheduled.Equal(cwf.Status.LastScheduledTime) {
15 ✔
186
                h.aw.Status.LastScheduled = cwf.Status.LastScheduledTime
2 ✔
187
                updated = true
2 ✔
188
        }
2 ✔
189
        if h.aw.Status.Failed != cwf.Status.Failed {
15 ✔
190
                h.aw.Status.Failed = cwf.Status.Failed
1 ✔
191
                updated = true
1 ✔
192
        }
1 ✔
193
        if h.aw.Status.Succeeded != cwf.Status.Succeeded {
15 ✔
194
                h.aw.Status.Succeeded = cwf.Status.Succeeded
1 ✔
195
                updated = true
1 ✔
196
        }
1 ✔
197

198
        // If the active workflow is not the same as the current one, update the reference
199
        if len(cwf.Status.Active) > 0 {
15 ✔
200
                // Should only contain a single element at most (expected to be in the same namespace!)
201
                ref := cwf.Status.Active[len(cwf.Status.Active)-1]
13 ✔
202

13 ✔
203
                if h.aw.Status.ActiveWorkflowRef.Name != ref.Name {
13 ✔
204
                        h.log.V(1).Info("Updating reference for cron workflow", "cronWorkflow", cwf.Name, "activeWorkflow", ref.Name)
8 ✔
205

8 ✔
206
                        // Get the active workflow
207
                        wf := wfv1alpha1.Workflow{}
8 ✔
208
                        if err := h.Get(ctx, namespacedName(h.aw.Namespace, ref.Name), &wf); err != nil {
8 ✔
209
                                return errLogAndWrap(h.log, err, "failed to fetch active workflow")
1 ✔
210
                        }
1 ✔
211

212
                        h.aw.Status.ActiveWorkflowRef = corev1.LocalObjectReference{
7 ✔
213
                                Name: wf.Name,
7 ✔
214
                        }
7 ✔
215
                        h.aw.Status.Message = ""
7 ✔
216
                        h.aw.Status.Phase = arcv1alpha1.WorkflowActive
7 ✔
217

7 ✔
218
                        changed, _ := h.setStatusFromWorkflow(ctx, h.log, h.aw, &wf)
7 ✔
219
                        updated = changed || updated
7 ✔
220
                }
221
        }
222

223
        // If there is an active workflow, check its status
224
        if h.aw.Status.ActiveWorkflowRef.Name != "" {
14 ✔
225
                wf := wfv1alpha1.Workflow{}
12 ✔
226
                if err := h.Get(ctx, namespacedName(h.aw.Namespace, h.aw.Status.ActiveWorkflowRef.Name), &wf); err != nil {
12 ✔
227
                        return errLogAndWrap(h.log, err, "failed to fetch active workflow")
×
228
                }
×
229

230
                changed, _ := h.setStatusFromWorkflow(ctx, h.log, h.aw, &wf)
12 ✔
231
                updated = changed || updated
12 ✔
232

12 ✔
233
                if wf.Status.Phase.Completed() {
12 ✔
234
                        h.aw.Status.ActiveWorkflowRef.Name = ""
7 ✔
235
                        updated = true
7 ✔
236
                }
7 ✔
237
        }
238

239
        if !updated {
14 ✔
240
                return nil
4 ✔
241
        }
4 ✔
242

243
        h.log.V(1).Info("Updating status from active workflow", "cronWorkflow", cwf.Name)
10 ✔
244

10 ✔
245
        if err := h.Status().Update(ctx, h.aw); err != nil {
10 ✔
246
                return errLogAndWrap(h.log, err, "failed to update status")
×
247
        }
×
248

249
        return nil
10 ✔
250
}
251

252
func hydrateArgoWorkflowSpec(aw *arcv1alpha1.ArtifactWorkflow, srcSecret *corev1.Secret, dstSecret *corev1.Secret) wfv1alpha1.WorkflowSpec {
253
        srcVolume := corev1.Volume{
2,132 ✔
254
                Name:     "src-secret-vol",
2,132 ✔
255
                EmptyDir: &corev1.EmptyDirVolumeSource{},
2,132 ✔
256
        }
2,132 ✔
257
        if srcSecret.Name != "" {
2,132 ✔
258
                srcVolume.VolumeSource = corev1.VolumeSource{
1,677 ✔
259
                        Secret: &corev1.SecretVolumeSource{
1,677 ✔
260
                                SecretName: srcSecret.Name,
1,677 ✔
261
                        },
1,677 ✔
262
                }
1,677 ✔
263
        }
264

265
        dstVolume := corev1.Volume{
2,132 ✔
266
                Name:     "dst-secret-vol",
2,132 ✔
267
                EmptyDir: &corev1.EmptyDirVolumeSource{},
2,132 ✔
268
        }
2,132 ✔
269
        if dstSecret.Name != "" {
2,132 ✔
270
                dstVolume.VolumeSource = corev1.VolumeSource{
1,677 ✔
271
                        Secret: &corev1.SecretVolumeSource{
1,677 ✔
272
                                SecretName: dstSecret.Name,
1,677 ✔
273
                        },
1,677 ✔
274
                }
1,677 ✔
275
        }
276

277
        parameters := []wfv1alpha1.Parameter{}
2,132 ✔
278
        for _, p := range aw.Spec.Parameters {
2,132 ✔
279
                parameters = append(parameters, wfv1alpha1.Parameter{
14,162 ✔
280
                        Name:  p.Name,
14,162 ✔
281
                        Value: (*wfv1alpha1.AnyString)(&p.Value),
14,162 ✔
282
                })
14,162 ✔
283
        }
14,162 ✔
284

285
        return wfv1alpha1.WorkflowSpec{
2,132 ✔
286
                WorkflowTemplateRef: &wfv1alpha1.WorkflowTemplateRef{
2,132 ✔
287
                        Name:         aw.Spec.WorkflowTemplateRef.Name,
2,132 ✔
288
                        ClusterScope: aw.Spec.WorkflowTemplateRef.ClusterScope,
2,132 ✔
289
                },
2,132 ✔
290
                Volumes: []corev1.Volume{
2,132 ✔
291
                        srcVolume,
2,132 ✔
292
                        dstVolume,
2,132 ✔
293
                },
2,132 ✔
294
                Arguments: wfv1alpha1.Arguments{
2,132 ✔
295
                        Parameters: parameters,
2,132 ✔
296
                },
2,132 ✔
297
        }
2,132 ✔
298
}
299

300
func hydrateArgoWorkflow(aw *arcv1alpha1.ArtifactWorkflow, srcSecret *corev1.Secret, dstSecret *corev1.Secret) *wfv1alpha1.Workflow {
301
        return &wfv1alpha1.Workflow{
2,118 ✔
302
                ObjectMeta: workflowObjectMeta(aw),
2,118 ✔
303
                Spec:       hydrateArgoWorkflowSpec(aw, srcSecret, dstSecret),
2,118 ✔
304
        }
2,118 ✔
305
}
306

307
func hydrateArgoCronWorkflow(aw *arcv1alpha1.ArtifactWorkflow, srcSecret *corev1.Secret, dstSecret *corev1.Secret) *wfv1alpha1.CronWorkflow {
308
        om := workflowObjectMeta(aw)
14 ✔
309
        wf := &wfv1alpha1.CronWorkflow{
14 ✔
310
                ObjectMeta: om,
14 ✔
311
                Spec: wfv1alpha1.CronWorkflowSpec{
14 ✔
312
                        WorkflowSpec:               hydrateArgoWorkflowSpec(aw, srcSecret, dstSecret),
14 ✔
313
                        Schedules:                  aw.Spec.Cron.Schedules,
14 ✔
314
                        ConcurrencyPolicy:          wfv1alpha1.ReplaceConcurrent,
14 ✔
315
                        StartingDeadlineSeconds:    aw.Spec.Cron.StartingDeadlineSeconds,
14 ✔
316
                        Timezone:                   aw.Spec.Cron.Timezone,
14 ✔
317
                        When:                       aw.Spec.Cron.When,
14 ✔
318
                        SuccessfulJobsHistoryLimit: new(int32(1)),
14 ✔
319
                        FailedJobsHistoryLimit:     new(int32(1)),
14 ✔
320
                        WorkflowMetadata:           &om,
14 ✔
321
                },
14 ✔
322
        }
14 ✔
323

324
        return wf
14 ✔
325
}
14 ✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc