2019-01-12 04:58:27 +00:00
|
|
|
/*
|
|
|
|
Copyright 2015 The Kubernetes Authors.
|
|
|
|
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
|
|
you may not use this file except in compliance with the License.
|
|
|
|
You may obtain a copy of the License at
|
|
|
|
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
See the License for the specific language governing permissions and
|
|
|
|
limitations under the License.
|
|
|
|
*/
|
|
|
|
|
|
|
|
package job
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
|
|
|
"strconv"
|
|
|
|
|
2019-04-07 17:07:55 +00:00
|
|
|
batchv1 "k8s.io/api/batch/v1"
|
2021-07-02 08:43:15 +00:00
|
|
|
apiequality "k8s.io/apimachinery/pkg/api/equality"
|
2019-01-12 04:58:27 +00:00
|
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
|
|
"k8s.io/apimachinery/pkg/fields"
|
|
|
|
"k8s.io/apimachinery/pkg/labels"
|
|
|
|
"k8s.io/apimachinery/pkg/runtime"
|
2019-04-07 17:07:55 +00:00
|
|
|
"k8s.io/apimachinery/pkg/runtime/schema"
|
2019-01-12 04:58:27 +00:00
|
|
|
"k8s.io/apimachinery/pkg/util/validation/field"
|
2019-04-07 17:07:55 +00:00
|
|
|
genericapirequest "k8s.io/apiserver/pkg/endpoints/request"
|
2019-01-12 04:58:27 +00:00
|
|
|
"k8s.io/apiserver/pkg/registry/generic"
|
|
|
|
"k8s.io/apiserver/pkg/registry/rest"
|
|
|
|
"k8s.io/apiserver/pkg/storage"
|
|
|
|
"k8s.io/apiserver/pkg/storage/names"
|
2019-08-30 18:33:25 +00:00
|
|
|
utilfeature "k8s.io/apiserver/pkg/util/feature"
|
2019-01-12 04:58:27 +00:00
|
|
|
"k8s.io/kubernetes/pkg/api/legacyscheme"
|
|
|
|
"k8s.io/kubernetes/pkg/api/pod"
|
|
|
|
"k8s.io/kubernetes/pkg/apis/batch"
|
|
|
|
"k8s.io/kubernetes/pkg/apis/batch/validation"
|
2021-07-02 08:43:15 +00:00
|
|
|
"k8s.io/kubernetes/pkg/apis/core"
|
2019-08-30 18:33:25 +00:00
|
|
|
"k8s.io/kubernetes/pkg/features"
|
2021-03-18 22:40:29 +00:00
|
|
|
"sigs.k8s.io/structured-merge-diff/v4/fieldpath"
|
2019-01-12 04:58:27 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
// jobStrategy implements verification logic for Replication Controllers.
|
|
|
|
type jobStrategy struct {
|
|
|
|
runtime.ObjectTyper
|
|
|
|
names.NameGenerator
|
|
|
|
}
|
|
|
|
|
|
|
|
// Strategy is the default logic that applies when creating and updating Replication Controller objects.
|
|
|
|
var Strategy = jobStrategy{legacyscheme.Scheme, names.SimpleNameGenerator}
|
|
|
|
|
2019-04-07 17:07:55 +00:00
|
|
|
// DefaultGarbageCollectionPolicy returns OrphanDependents for batch/v1 for backwards compatibility,
|
|
|
|
// and DeleteDependents for all other versions.
|
2019-01-12 04:58:27 +00:00
|
|
|
func (jobStrategy) DefaultGarbageCollectionPolicy(ctx context.Context) rest.GarbageCollectionPolicy {
|
2019-04-07 17:07:55 +00:00
|
|
|
var groupVersion schema.GroupVersion
|
|
|
|
if requestInfo, found := genericapirequest.RequestInfoFrom(ctx); found {
|
|
|
|
groupVersion = schema.GroupVersion{Group: requestInfo.APIGroup, Version: requestInfo.APIVersion}
|
|
|
|
}
|
|
|
|
switch groupVersion {
|
|
|
|
case batchv1.SchemeGroupVersion:
|
|
|
|
// for back compatibility
|
|
|
|
return rest.OrphanDependents
|
|
|
|
default:
|
|
|
|
return rest.DeleteDependents
|
|
|
|
}
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// NamespaceScoped returns true because all jobs need to be within a namespace.
|
|
|
|
func (jobStrategy) NamespaceScoped() bool {
|
|
|
|
return true
|
|
|
|
}
|
|
|
|
|
2021-03-18 22:40:29 +00:00
|
|
|
// GetResetFields returns the set of fields that get reset by the strategy
|
|
|
|
// and should not be modified by the user.
|
|
|
|
func (jobStrategy) GetResetFields() map[fieldpath.APIVersion]*fieldpath.Set {
|
|
|
|
fields := map[fieldpath.APIVersion]*fieldpath.Set{
|
|
|
|
"batch/v1": fieldpath.NewSet(
|
|
|
|
fieldpath.MakePathOrDie("status"),
|
|
|
|
),
|
|
|
|
}
|
|
|
|
|
|
|
|
return fields
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
// PrepareForCreate clears the status of a job before creation.
|
|
|
|
func (jobStrategy) PrepareForCreate(ctx context.Context, obj runtime.Object) {
|
|
|
|
job := obj.(*batch.Job)
|
|
|
|
job.Status = batch.JobStatus{}
|
2019-08-30 18:33:25 +00:00
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
job.Generation = 1
|
|
|
|
|
2019-08-30 18:33:25 +00:00
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.TTLAfterFinished) {
|
|
|
|
job.Spec.TTLSecondsAfterFinished = nil
|
|
|
|
}
|
2019-01-12 04:58:27 +00:00
|
|
|
|
2021-03-18 22:40:29 +00:00
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.IndexedJob) {
|
|
|
|
job.Spec.CompletionMode = nil
|
|
|
|
}
|
|
|
|
|
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.SuspendJob) {
|
|
|
|
job.Spec.Suspend = nil
|
|
|
|
}
|
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
if utilfeature.DefaultFeatureGate.Enabled(features.JobTrackingWithFinalizers) {
|
|
|
|
// Until this feature graduates to GA and soaks in clusters, we use an
|
|
|
|
// annotation to mark whether jobs are tracked with it.
|
|
|
|
addJobTrackingAnnotation(job)
|
|
|
|
} else {
|
|
|
|
dropJobTrackingAnnotation(job)
|
|
|
|
}
|
|
|
|
|
2019-04-07 17:07:55 +00:00
|
|
|
pod.DropDisabledTemplateFields(&job.Spec.Template, nil)
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
func addJobTrackingAnnotation(job *batch.Job) {
|
|
|
|
if job.Annotations == nil {
|
|
|
|
job.Annotations = map[string]string{}
|
|
|
|
}
|
|
|
|
job.Annotations[batchv1.JobTrackingFinalizer] = ""
|
|
|
|
}
|
|
|
|
|
|
|
|
func hasJobTrackingAnnotation(job *batch.Job) bool {
|
|
|
|
if job.Annotations == nil {
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
_, ok := job.Annotations[batchv1.JobTrackingFinalizer]
|
|
|
|
return ok
|
|
|
|
}
|
|
|
|
|
|
|
|
func dropJobTrackingAnnotation(job *batch.Job) {
|
|
|
|
delete(job.Annotations, batchv1.JobTrackingFinalizer)
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
// PrepareForUpdate clears fields that are not allowed to be set by end users on update.
|
|
|
|
func (jobStrategy) PrepareForUpdate(ctx context.Context, obj, old runtime.Object) {
|
|
|
|
newJob := obj.(*batch.Job)
|
|
|
|
oldJob := old.(*batch.Job)
|
|
|
|
newJob.Status = oldJob.Status
|
|
|
|
|
2019-08-30 18:33:25 +00:00
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.TTLAfterFinished) && oldJob.Spec.TTLSecondsAfterFinished == nil {
|
2019-04-07 17:07:55 +00:00
|
|
|
newJob.Spec.TTLSecondsAfterFinished = nil
|
|
|
|
}
|
2019-01-12 04:58:27 +00:00
|
|
|
|
2021-03-18 22:40:29 +00:00
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.IndexedJob) && oldJob.Spec.CompletionMode == nil {
|
|
|
|
newJob.Spec.CompletionMode = nil
|
|
|
|
}
|
|
|
|
|
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.SuspendJob) {
|
|
|
|
// There are 3 possible values (nil, true, false) for each flag, so 9
|
|
|
|
// combinations. We want to disallow everything except true->false and
|
|
|
|
// true->nil when the feature gate is disabled. Or, basically allow this
|
|
|
|
// only when oldJob is true.
|
|
|
|
if oldJob.Spec.Suspend == nil || !*oldJob.Spec.Suspend {
|
|
|
|
newJob.Spec.Suspend = oldJob.Spec.Suspend
|
|
|
|
}
|
|
|
|
}
|
2021-07-02 08:43:15 +00:00
|
|
|
if !utilfeature.DefaultFeatureGate.Enabled(features.JobTrackingWithFinalizers) && !hasJobTrackingAnnotation(oldJob) {
|
|
|
|
dropJobTrackingAnnotation(newJob)
|
|
|
|
}
|
2021-03-18 22:40:29 +00:00
|
|
|
|
2019-04-07 17:07:55 +00:00
|
|
|
pod.DropDisabledTemplateFields(&newJob.Spec.Template, &oldJob.Spec.Template)
|
2021-07-02 08:43:15 +00:00
|
|
|
|
|
|
|
// Any changes to the spec increment the generation number.
|
|
|
|
// See metav1.ObjectMeta description for more information on Generation.
|
|
|
|
if !apiequality.Semantic.DeepEqual(newJob.Spec, oldJob.Spec) {
|
|
|
|
newJob.Generation = oldJob.Generation + 1
|
|
|
|
}
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Validate validates a new job.
|
|
|
|
func (jobStrategy) Validate(ctx context.Context, obj runtime.Object) field.ErrorList {
|
|
|
|
job := obj.(*batch.Job)
|
|
|
|
// TODO: move UID generation earlier and do this in defaulting logic?
|
|
|
|
if job.Spec.ManualSelector == nil || *job.Spec.ManualSelector == false {
|
|
|
|
generateSelector(job)
|
|
|
|
}
|
2021-07-02 08:43:15 +00:00
|
|
|
opts := validationOptionsForJob(job, nil)
|
2020-12-01 01:06:26 +00:00
|
|
|
return validation.ValidateJob(job, opts)
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
func validationOptionsForJob(newJob, oldJob *batch.Job) validation.JobValidationOptions {
|
|
|
|
var newPodTemplate, oldPodTemplate *core.PodTemplateSpec
|
|
|
|
if newJob != nil {
|
|
|
|
newPodTemplate = &newJob.Spec.Template
|
|
|
|
}
|
|
|
|
if oldJob != nil {
|
|
|
|
oldPodTemplate = &oldJob.Spec.Template
|
|
|
|
}
|
|
|
|
opts := validation.JobValidationOptions{
|
|
|
|
PodValidationOptions: pod.GetValidationOptionsFromPodTemplate(newPodTemplate, oldPodTemplate),
|
|
|
|
AllowTrackingAnnotation: utilfeature.DefaultFeatureGate.Enabled(features.JobTrackingWithFinalizers),
|
|
|
|
}
|
|
|
|
if oldJob != nil {
|
|
|
|
// Because we don't support the tracking with finalizers for already
|
|
|
|
// existing jobs, we allow the annotation only if the Job already had it,
|
|
|
|
// regardless of the feature gate.
|
|
|
|
opts.AllowTrackingAnnotation = hasJobTrackingAnnotation(oldJob)
|
|
|
|
}
|
|
|
|
return opts
|
|
|
|
}
|
|
|
|
|
|
|
|
// WarningsOnCreate returns warnings for the creation of the given object.
|
|
|
|
func (jobStrategy) WarningsOnCreate(ctx context.Context, obj runtime.Object) []string {
|
|
|
|
newJob := obj.(*batch.Job)
|
|
|
|
return pod.GetWarningsForPodTemplate(ctx, field.NewPath("spec", "template"), &newJob.Spec.Template, nil)
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
// generateSelector adds a selector to a job and labels to its template
|
|
|
|
// which can be used to uniquely identify the pods created by that job,
|
|
|
|
// if the user has requested this behavior.
|
|
|
|
func generateSelector(obj *batch.Job) {
|
|
|
|
if obj.Spec.Template.Labels == nil {
|
|
|
|
obj.Spec.Template.Labels = make(map[string]string)
|
|
|
|
}
|
|
|
|
// The job-name label is unique except in cases that are expected to be
|
|
|
|
// quite uncommon, and is more user friendly than uid. So, we add it as
|
|
|
|
// a label.
|
|
|
|
_, found := obj.Spec.Template.Labels["job-name"]
|
|
|
|
if found {
|
|
|
|
// User asked us to not automatically generate a selector and labels,
|
|
|
|
// but set a possibly conflicting value. If there is a conflict,
|
|
|
|
// we will reject in validation.
|
|
|
|
} else {
|
|
|
|
obj.Spec.Template.Labels["job-name"] = string(obj.ObjectMeta.Name)
|
|
|
|
}
|
|
|
|
// The controller-uid label makes the pods that belong to this job
|
|
|
|
// only match this job.
|
|
|
|
_, found = obj.Spec.Template.Labels["controller-uid"]
|
|
|
|
if found {
|
|
|
|
// User asked us to automatically generate a selector and labels,
|
|
|
|
// but set a possibly conflicting value. If there is a conflict,
|
|
|
|
// we will reject in validation.
|
|
|
|
} else {
|
|
|
|
obj.Spec.Template.Labels["controller-uid"] = string(obj.ObjectMeta.UID)
|
|
|
|
}
|
|
|
|
// Select the controller-uid label. This is sufficient for uniqueness.
|
|
|
|
if obj.Spec.Selector == nil {
|
|
|
|
obj.Spec.Selector = &metav1.LabelSelector{}
|
|
|
|
}
|
|
|
|
if obj.Spec.Selector.MatchLabels == nil {
|
|
|
|
obj.Spec.Selector.MatchLabels = make(map[string]string)
|
|
|
|
}
|
|
|
|
if _, found := obj.Spec.Selector.MatchLabels["controller-uid"]; !found {
|
|
|
|
obj.Spec.Selector.MatchLabels["controller-uid"] = string(obj.ObjectMeta.UID)
|
|
|
|
}
|
|
|
|
// If the user specified matchLabel controller-uid=$WRONGUID, then it should fail
|
|
|
|
// in validation, either because the selector does not match the pod template
|
|
|
|
// (controller-uid=$WRONGUID does not match controller-uid=$UID, which we applied
|
|
|
|
// above, or we will reject in validation because the template has the wrong
|
|
|
|
// labels.
|
|
|
|
}
|
|
|
|
|
|
|
|
// TODO: generalize generateSelector so it can work for other controller
|
|
|
|
// objects such as ReplicaSet. Can use pkg/api/meta to generically get the
|
|
|
|
// UID, but need some way to generically access the selector and pod labels
|
|
|
|
// fields.
|
|
|
|
|
|
|
|
// Canonicalize normalizes the object after validation.
|
|
|
|
func (jobStrategy) Canonicalize(obj runtime.Object) {
|
|
|
|
}
|
|
|
|
|
|
|
|
func (jobStrategy) AllowUnconditionalUpdate() bool {
|
|
|
|
return true
|
|
|
|
}
|
|
|
|
|
|
|
|
// AllowCreateOnUpdate is false for jobs; this means a POST is needed to create one.
|
|
|
|
func (jobStrategy) AllowCreateOnUpdate() bool {
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
|
|
|
|
// ValidateUpdate is the default update validation for an end user.
|
|
|
|
func (jobStrategy) ValidateUpdate(ctx context.Context, obj, old runtime.Object) field.ErrorList {
|
2019-04-07 17:07:55 +00:00
|
|
|
job := obj.(*batch.Job)
|
|
|
|
oldJob := old.(*batch.Job)
|
2020-12-01 01:06:26 +00:00
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
opts := validationOptionsForJob(job, oldJob)
|
2020-12-01 01:06:26 +00:00
|
|
|
validationErrorList := validation.ValidateJob(job, opts)
|
2021-07-02 08:43:15 +00:00
|
|
|
updateErrorList := validation.ValidateJobUpdate(job, oldJob, opts.PodValidationOptions)
|
2019-01-12 04:58:27 +00:00
|
|
|
return append(validationErrorList, updateErrorList...)
|
|
|
|
}
|
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
// WarningsOnUpdate returns warnings for the given update.
|
|
|
|
func (jobStrategy) WarningsOnUpdate(ctx context.Context, obj, old runtime.Object) []string {
|
|
|
|
var warnings []string
|
|
|
|
newJob := obj.(*batch.Job)
|
|
|
|
oldJob := old.(*batch.Job)
|
|
|
|
if newJob.Generation != oldJob.Generation {
|
|
|
|
warnings = pod.GetWarningsForPodTemplate(ctx, field.NewPath("spec", "template"), &newJob.Spec.Template, &oldJob.Spec.Template)
|
|
|
|
}
|
|
|
|
return warnings
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
type jobStatusStrategy struct {
|
|
|
|
jobStrategy
|
|
|
|
}
|
|
|
|
|
|
|
|
var StatusStrategy = jobStatusStrategy{Strategy}
|
|
|
|
|
2021-03-18 22:40:29 +00:00
|
|
|
// GetResetFields returns the set of fields that get reset by the strategy
|
|
|
|
// and should not be modified by the user.
|
|
|
|
func (jobStatusStrategy) GetResetFields() map[fieldpath.APIVersion]*fieldpath.Set {
|
|
|
|
return map[fieldpath.APIVersion]*fieldpath.Set{
|
|
|
|
"batch/v1": fieldpath.NewSet(
|
|
|
|
fieldpath.MakePathOrDie("spec"),
|
|
|
|
),
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
func (jobStatusStrategy) PrepareForUpdate(ctx context.Context, obj, old runtime.Object) {
|
|
|
|
newJob := obj.(*batch.Job)
|
|
|
|
oldJob := old.(*batch.Job)
|
|
|
|
newJob.Spec = oldJob.Spec
|
|
|
|
}
|
|
|
|
|
|
|
|
func (jobStatusStrategy) ValidateUpdate(ctx context.Context, obj, old runtime.Object) field.ErrorList {
|
|
|
|
return validation.ValidateJobUpdateStatus(obj.(*batch.Job), old.(*batch.Job))
|
|
|
|
}
|
|
|
|
|
2021-07-02 08:43:15 +00:00
|
|
|
// WarningsOnUpdate returns warnings for the given update.
|
|
|
|
func (jobStatusStrategy) WarningsOnUpdate(ctx context.Context, obj, old runtime.Object) []string {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2019-01-12 04:58:27 +00:00
|
|
|
// JobSelectableFields returns a field set that represents the object for matching purposes.
|
|
|
|
func JobToSelectableFields(job *batch.Job) fields.Set {
|
|
|
|
objectMetaFieldsSet := generic.ObjectMetaFieldsSet(&job.ObjectMeta, true)
|
|
|
|
specificFieldsSet := fields.Set{
|
|
|
|
"status.successful": strconv.Itoa(int(job.Status.Succeeded)),
|
|
|
|
}
|
|
|
|
return generic.MergeFieldsSets(objectMetaFieldsSet, specificFieldsSet)
|
|
|
|
}
|
|
|
|
|
|
|
|
// GetAttrs returns labels and fields of a given object for filtering purposes.
|
2019-04-07 17:07:55 +00:00
|
|
|
func GetAttrs(obj runtime.Object) (labels.Set, fields.Set, error) {
|
2019-01-12 04:58:27 +00:00
|
|
|
job, ok := obj.(*batch.Job)
|
|
|
|
if !ok {
|
2019-04-07 17:07:55 +00:00
|
|
|
return nil, nil, fmt.Errorf("given object is not a job.")
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
2019-04-07 17:07:55 +00:00
|
|
|
return labels.Set(job.ObjectMeta.Labels), JobToSelectableFields(job), nil
|
2019-01-12 04:58:27 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// MatchJob is the filter used by the generic etcd backend to route
|
|
|
|
// watch events from etcd to clients of the apiserver only interested in specific
|
|
|
|
// labels/fields.
|
|
|
|
func MatchJob(label labels.Selector, field fields.Selector) storage.SelectionPredicate {
|
|
|
|
return storage.SelectionPredicate{
|
|
|
|
Label: label,
|
|
|
|
Field: field,
|
|
|
|
GetAttrs: GetAttrs,
|
|
|
|
}
|
|
|
|
}
|