From 86effb72a847e82a296105d91a8c2f60f36ef8c1 Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Fri, 18 Sep 2026 18:20:10 +0000 Subject: [PATCH 1/3] feat: execute jobs on Kubernetes --- cmd/main.go | 8 + config/rbac/role.yaml | 59 ++- internal/adapter/kubernetes/job.go | 132 +++++++ internal/adapter/kubernetes/job_test.go | 71 ++++ internal/controller/job_controller.go | 397 +++++++++++++++++++++ internal/controller/job_controller_test.go | 239 +++++++++++++ 6 files changed, 900 insertions(+), 6 deletions(-) create mode 100644 internal/adapter/kubernetes/job.go create mode 100644 internal/adapter/kubernetes/job_test.go create mode 100644 internal/controller/job_controller.go create mode 100644 internal/controller/job_controller_test.go diff --git a/cmd/main.go b/cmd/main.go index d5bb38c..f7e62d4 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -20,6 +20,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/webhook" executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + "git.ddupan.top/panxiao81/ayatori/internal/controller" // +kubebuilder:scaffold:imports ) @@ -165,6 +166,13 @@ func main() { os.Exit(1) } + if err := (&controller.JobReconciler{ + Client: mgr.GetClient(), + }).SetupWithManager(mgr); err != nil { + setupLog.Error(err, "Failed to create controller", "controller", "Job") + os.Exit(1) + } + // +kubebuilder:scaffold:builder if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil { diff --git a/config/rbac/role.yaml b/config/rbac/role.yaml index 759ce82..61c28db 100644 --- a/config/rbac/role.yaml +++ b/config/rbac/role.yaml @@ -1,11 +1,58 @@ +--- apiVersion: rbac.authorization.k8s.io/v1 kind: ClusterRole metadata: - labels: - app.kubernetes.io/name: ayatori - app.kubernetes.io/managed-by: kustomize name: manager-role rules: -- apiGroups: [""] - resources: ["pods"] - verbs: ["get", "list", "watch"] +- apiGroups: + - "" + resources: + - namespaces + - serviceaccounts + verbs: + - get + - list + - watch +- apiGroups: + - batch + resources: + - jobs + verbs: + - create + - delete + - get + - list + - watch +- apiGroups: + - execution.ayatori.ddupan.top + resources: + - jobclasses + - kubernetesexecutionparameters + verbs: + - get + - list + - watch +- apiGroups: + - execution.ayatori.ddupan.top + resources: + - jobs + verbs: + - get + - list + - patch + - update + - watch +- apiGroups: + - execution.ayatori.ddupan.top + resources: + - jobs/finalizers + verbs: + - update +- apiGroups: + - execution.ayatori.ddupan.top + resources: + - jobs/status + verbs: + - get + - patch + - update diff --git a/internal/adapter/kubernetes/job.go b/internal/adapter/kubernetes/job.go new file mode 100644 index 0000000..5d06f49 --- /dev/null +++ b/internal/adapter/kubernetes/job.go @@ -0,0 +1,132 @@ +package kubernetes + +import ( + "fmt" + + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +const ( + ControllerName = "execution.ayatori.ddupan.top/kubernetes" + ReferenceType = "Job" + JobUIDLabel = "execution.ayatori.ddupan.top/job-uid" +) + +var ayatoriJobGVK = schema.GroupVersionKind{ + Group: executionv1alpha1.GroupVersion.Group, + Version: executionv1alpha1.GroupVersion.Version, + Kind: "Job", +} + +// BuildJob translates the stable execution API into the Kubernetes adapter's +// backend object. It intentionally does not accept or expose a PodSpec. +func BuildJob( + job *executionv1alpha1.Job, + parameters *executionv1alpha1.KubernetesExecutionParameters, + resources executionv1alpha1.ExecutionResourceRequirements, +) *batchv1.Job { + backoffLimit := int32(0) + controller := true + blockOwnerDeletion := true + + //nolint:modernize // ObjectMeta is promoted through embedded TypeMeta; embedlit produces invalid Go here. + return &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{ + Name: job.Name, + Namespace: job.Namespace, + Labels: map[string]string{ + JobUIDLabel: string(job.UID), + }, + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: ayatoriJobGVK.GroupVersion().String(), + Kind: ayatoriJobGVK.Kind, + Name: job.Name, + UID: job.UID, + Controller: &controller, + BlockOwnerDeletion: &blockOwnerDeletion, + }}, + }, + Spec: batchv1.JobSpec{ + BackoffLimit: &backoffLimit, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{JobUIDLabel: string(job.UID)}}, + Spec: corev1.PodSpec{ + RestartPolicy: corev1.RestartPolicyNever, + ServiceAccountName: parameters.Spec.ServiceAccountName, + RuntimeClassName: optionalString(parameters.Spec.RuntimeClassName), + NodeSelector: parameters.Spec.Scheduling.NodeSelector, + Tolerations: parameters.Spec.Scheduling.Tolerations, + SecurityContext: parameters.Spec.PodSecurityContext, + ImagePullSecrets: job.Spec.Task.ImagePullSecrets, + Containers: []corev1.Container{{ + Name: "task", + Image: job.Spec.Task.Image, + ImagePullPolicy: parameters.Spec.ImagePullPolicy, + Command: job.Spec.Task.Command, + Args: job.Spec.Task.Args, + WorkingDir: job.Spec.Task.WorkingDir, + Env: environment(job.Spec.Task.Env), + Resources: resourceRequirements(resources), + }}, + }, + }, + }, + } +} + +func ValidateOwnership(owner *executionv1alpha1.Job, backend *batchv1.Job) error { + if backend.Labels[JobUIDLabel] != string(owner.UID) { + return fmt.Errorf("backend Job %s/%s is not owned by Ayatori Job UID %s", backend.Namespace, backend.Name, owner.UID) + } + return nil +} + +func environment(values []executionv1alpha1.EnvVar) []corev1.EnvVar { + result := make([]corev1.EnvVar, 0, len(values)) + for _, value := range values { + env := corev1.EnvVar{Name: value.Name} + if value.Value != nil { + env.Value = *value.Value + } + if value.ValueFrom != nil { + env.ValueFrom = &corev1.EnvVarSource{ + SecretKeyRef: value.ValueFrom.SecretKeyRef, + ConfigMapKeyRef: value.ValueFrom.ConfigMapKeyRef, + } + } + result = append(result, env) + } + return result +} + +func resourceRequirements(resources executionv1alpha1.ExecutionResourceRequirements) corev1.ResourceRequirements { + return corev1.ResourceRequirements{ + Requests: resourceList(resources.Requests), + Limits: resourceList(resources.Limits), + } +} + +func resourceList(values executionv1alpha1.ResourceValues) corev1.ResourceList { + result := corev1.ResourceList{} + if values.CPU != nil { + result[corev1.ResourceCPU] = values.CPU.DeepCopy() + } + if values.Memory != nil { + result[corev1.ResourceMemory] = values.Memory.DeepCopy() + } + if len(result) == 0 { + return nil + } + return result +} + +func optionalString(value string) *string { + if value == "" { + return nil + } + return &value +} diff --git a/internal/adapter/kubernetes/job_test.go b/internal/adapter/kubernetes/job_test.go new file mode 100644 index 0000000..4887329 --- /dev/null +++ b/internal/adapter/kubernetes/job_test.go @@ -0,0 +1,71 @@ +package kubernetes + +import ( + "testing" + + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" +) + +func TestBuildJob(t *testing.T) { + literal := "world" + cpuRequest := resource.MustParse("100m") + memoryLimit := resource.MustParse("128Mi") + //nolint:modernize // ObjectMeta is promoted through embedded TypeMeta; embedlit produces invalid Go here. + job := &executionv1alpha1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "hello", Namespace: "ci", UID: types.UID("job-uid")}, + Spec: executionv1alpha1.JobSpec{Task: executionv1alpha1.TaskSpec{ + Image: "alpine:3.22", Command: []string{"echo"}, Args: []string{"hello"}, + Env: []executionv1alpha1.EnvVar{ + {Name: "TARGET", Value: &literal}, + {Name: "TOKEN", ValueFrom: &executionv1alpha1.EnvVarSource{ + //nolint:modernize // LocalObjectReference is an embedded Kubernetes API field. + SecretKeyRef: &corev1.SecretKeySelector{LocalObjectReference: corev1.LocalObjectReference{Name: "token"}, Key: "value"}, + }}, + }, + }}, + } + parameters := &executionv1alpha1.KubernetesExecutionParameters{Spec: executionv1alpha1.KubernetesExecutionParametersSpec{ + ServiceAccountName: "runner", RuntimeClassName: "runc", ImagePullPolicy: corev1.PullIfNotPresent, + Scheduling: executionv1alpha1.KubernetesSchedulingParameters{NodeSelector: map[string]string{"role": "execution"}}, + }} + resources := executionv1alpha1.ExecutionResourceRequirements{ + Requests: executionv1alpha1.ResourceValues{CPU: &cpuRequest}, + Limits: executionv1alpha1.ResourceValues{Memory: &memoryLimit}, + } + + backend := BuildJob(job, parameters, resources) + pod := backend.Spec.Template.Spec + if backend.Spec.BackoffLimit == nil || *backend.Spec.BackoffLimit != 0 { + t.Fatalf("backoffLimit = %v, want 0", backend.Spec.BackoffLimit) + } + if pod.RestartPolicy != corev1.RestartPolicyNever || pod.ServiceAccountName != "runner" { + t.Fatalf("unexpected pod execution policy: %#v", pod) + } + if pod.RuntimeClassName == nil || *pod.RuntimeClassName != "runc" { + t.Fatalf("runtimeClassName = %v, want runc", pod.RuntimeClassName) + } + container := pod.Containers[0] + if container.Resources.Requests.Cpu().Cmp(cpuRequest) != 0 || container.Resources.Limits.Memory().Cmp(memoryLimit) != 0 { + t.Fatalf("resources were not mapped: %#v", container.Resources) + } + if container.Env[1].ValueFrom == nil || container.Env[1].ValueFrom.SecretKeyRef.Name != "token" { + t.Fatalf("secret reference was not preserved: %#v", container.Env[1]) + } + if backend.Labels[JobUIDLabel] != "job-uid" || backend.OwnerReferences[0].UID != job.UID { + t.Fatalf("ownership identity was not preserved: %#v", backend.ObjectMeta) + } +} + +func TestValidateOwnership(t *testing.T) { + job := &executionv1alpha1.Job{} + job.UID = types.UID("expected") + backend := BuildJob(job, &executionv1alpha1.KubernetesExecutionParameters{}, executionv1alpha1.ExecutionResourceRequirements{}) + backend.Labels[JobUIDLabel] = "different" + if err := ValidateOwnership(job, backend); err == nil { + t.Fatal("ValidateOwnership() succeeded for a different Job UID") + } +} diff --git a/internal/controller/job_controller.go b/internal/controller/job_controller.go new file mode 100644 index 0000000..7131110 --- /dev/null +++ b/internal/controller/job_controller.go @@ -0,0 +1,397 @@ +package controller + +import ( + "context" + "fmt" + "slices" + "time" + + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + kubernetesadapter "git.ddupan.top/panxiao81/ayatori/internal/adapter/kubernetes" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/log" +) + +const ( + jobFinalizer = "execution.ayatori.ddupan.top/job-cleanup" + reasonResultUnknown = "ResultUnknown" +) + +// JobReconciler executes Ayatori Jobs using supported adapters. +type JobReconciler struct { + client.Client + Now func() time.Time +} + +// +kubebuilder:rbac:groups=execution.ayatori.ddupan.top,resources=jobs,verbs=get;list;watch;update;patch +// +kubebuilder:rbac:groups=execution.ayatori.ddupan.top,resources=jobs/status,verbs=get;update;patch +// +kubebuilder:rbac:groups=execution.ayatori.ddupan.top,resources=jobs/finalizers,verbs=update +// +kubebuilder:rbac:groups=execution.ayatori.ddupan.top,resources=jobclasses;kubernetesexecutionparameters,verbs=get;list;watch +// +kubebuilder:rbac:groups=batch,resources=jobs,verbs=get;list;watch;create;delete +// +kubebuilder:rbac:groups="",resources=namespaces;serviceaccounts,verbs=get;list;watch + +func (r *JobReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) { + logger := log.FromContext(ctx) + job := &executionv1alpha1.Job{} + if err := r.Get(ctx, request.NamespacedName, job); err != nil { + return ctrl.Result{}, client.IgnoreNotFound(err) + } + + if !job.DeletionTimestamp.IsZero() { + return ctrl.Result{}, r.finalize(ctx, job) + } + if isTerminal(job) { + return ctrl.Result{}, nil + } + if job.Spec.DesiredState == executionv1alpha1.JobDesiredStateCancelled { + return ctrl.Result{}, r.cancel(ctx, job) + } + + if !containsString(job.Finalizers, jobFinalizer) { + job.Finalizers = append(job.Finalizers, jobFinalizer) + if err := r.Update(ctx, job); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{}, nil + } + if job.Status.Execution != nil { + return ctrl.Result{}, r.observeExisting(ctx, job) + } + + class, parameters, resources, waiting, err := r.resolve(ctx, job) + if err != nil { + return ctrl.Result{}, err + } + if waiting { + return ctrl.Result{RequeueAfter: 30 * time.Second}, nil + } + + backend := &batchv1.Job{} + key := types.NamespacedName{Namespace: job.Namespace, Name: job.Name} + err = r.Get(ctx, key, backend) + if apierrors.IsNotFound(err) { + backend = kubernetesadapter.BuildJob(job, parameters, resources) + if err := r.Create(ctx, backend); err != nil { + return ctrl.Result{}, err + } + logger.Info("Created Kubernetes backend Job", "backend", key) + return ctrl.Result{}, r.markScheduled(ctx, job, class, parameters, resources, backend) + } + if err != nil { + return ctrl.Result{}, err + } + if err := kubernetesadapter.ValidateOwnership(job, backend); err != nil { + return ctrl.Result{}, r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionScheduled, Status: metav1.ConditionFalse, + Reason: "BackendConflict", Message: err.Error(), + }) + } + if !conditionTrue(job.Status.Conditions, executionv1alpha1.JobConditionScheduled) { + return ctrl.Result{}, r.markScheduled(ctx, job, class, parameters, resources, backend) + } + return ctrl.Result{}, r.observe(ctx, job, backend) +} + +func (r *JobReconciler) observeExisting(ctx context.Context, job *executionv1alpha1.Job) error { + if job.Status.Execution.Adapter != "kubernetes" { + return r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionUnknown, + Reason: reasonResultUnknown, Message: fmt.Sprintf("adapter %q is not available", job.Status.Execution.Adapter), + }) + } + backend := &batchv1.Job{} + key := types.NamespacedName{Namespace: job.Namespace, Name: job.Name} + if err := r.Get(ctx, key, backend); err != nil { + if apierrors.IsNotFound(err) { + return r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionUnknown, + Reason: reasonResultUnknown, Message: "Kubernetes backend Job is missing", + }) + } + return err + } + if err := kubernetesadapter.ValidateOwnership(job, backend); err != nil { + return r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionUnknown, + Reason: reasonResultUnknown, Message: err.Error(), + }) + } + return r.observe(ctx, job, backend) +} + +func (r *JobReconciler) resolve( + ctx context.Context, + job *executionv1alpha1.Job, +) (*executionv1alpha1.JobClass, *executionv1alpha1.KubernetesExecutionParameters, executionv1alpha1.ExecutionResourceRequirements, bool, error) { + if job.Spec.JobClassName == "" { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "NoDefaultJobClass", "spec.jobClassName is required in the first implementation slice") + } + + class := &executionv1alpha1.JobClass{} + if err := r.Get(ctx, types.NamespacedName{Name: job.Spec.JobClassName}, class); err != nil { + if apierrors.IsNotFound(err) { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "JobClassNotFound", fmt.Sprintf("JobClass %q does not exist", job.Spec.JobClassName)) + } + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, false, err + } + if class.Spec.ControllerName != kubernetesadapter.ControllerName { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "UnsupportedController", fmt.Sprintf("controller %q is not supported", class.Spec.ControllerName)) + } + ref := class.Spec.ParametersRef + if ref.Group != executionv1alpha1.GroupVersion.Group || ref.Kind != "KubernetesExecutionParameters" { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "InvalidParametersReference", "JobClass must reference KubernetesExecutionParameters") + } + if allowed, err := r.namespaceAllowed(ctx, job.Namespace, class.Spec.AllowedNamespaces); err != nil { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, false, err + } else if !allowed { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "NamespaceNotAllowed", fmt.Sprintf("namespace %q is not allowed by JobClass %q", job.Namespace, class.Name)) + } + + parameters := &executionv1alpha1.KubernetesExecutionParameters{} + if err := r.Get(ctx, types.NamespacedName{Name: ref.Name}, parameters); err != nil { + if apierrors.IsNotFound(err) { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "ParametersNotFound", fmt.Sprintf("KubernetesExecutionParameters %q does not exist", ref.Name)) + } + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, false, err + } + serviceAccount := &corev1.ServiceAccount{} + if err := r.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: parameters.Spec.ServiceAccountName}, serviceAccount); err != nil { + if apierrors.IsNotFound(err) { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "ServiceAccountNotFound", fmt.Sprintf("ServiceAccount %q does not exist", parameters.Spec.ServiceAccountName)) + } + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, false, err + } + + resources := applyResourceDefaults(job.Spec.Resources, class.Spec.Resources.Defaults) + if err := validateResources(resources); err != nil { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, true, r.reject(ctx, job, "InvalidResources", err.Error()) + } + if err := r.accept(ctx, job, class, parameters, resources); err != nil { + return nil, nil, executionv1alpha1.ExecutionResourceRequirements{}, false, err + } + return class, parameters, resources, false, nil +} + +func (r *JobReconciler) namespaceAllowed(ctx context.Context, namespace string, selector *metav1.LabelSelector) (bool, error) { + if selector == nil { + return true, nil + } + ns := &corev1.Namespace{} + if err := r.Get(ctx, types.NamespacedName{Name: namespace}, ns); err != nil { + return false, err + } + compiled, err := metav1.LabelSelectorAsSelector(selector) + if err != nil { + return false, err + } + return compiled.Matches(labels.Set(ns.Labels)), nil +} + +func (r *JobReconciler) accept(ctx context.Context, job *executionv1alpha1.Job, class *executionv1alpha1.JobClass, parameters *executionv1alpha1.KubernetesExecutionParameters, resources executionv1alpha1.ExecutionResourceRequirements) error { + job.Status.ResolvedJobClass = &executionv1alpha1.ResolvedJobClassReference{ + Name: class.Name, UID: class.UID, ControllerName: class.Spec.ControllerName, + ParametersRef: executionv1alpha1.ParametersReference{ + Group: class.Spec.ParametersRef.Group, Kind: class.Spec.ParametersRef.Kind, + Name: parameters.Name, UID: parameters.UID, + }, + } + job.Status.EffectiveResources = resources + return r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionAccepted, Status: metav1.ConditionTrue, + Reason: "Accepted", Message: fmt.Sprintf("JobClass %q accepted", class.Name), + }) +} + +func (r *JobReconciler) reject(ctx context.Context, job *executionv1alpha1.Job, reason, message string) error { + return r.setCondition(ctx, job, metav1.Condition{ + Type: executionv1alpha1.JobConditionAccepted, Status: metav1.ConditionFalse, + Reason: reason, Message: message, + }) +} + +func (r *JobReconciler) markScheduled(ctx context.Context, job *executionv1alpha1.Job, class *executionv1alpha1.JobClass, parameters *executionv1alpha1.KubernetesExecutionParameters, resources executionv1alpha1.ExecutionResourceRequirements, backend *batchv1.Job) error { + job.Status.ResolvedJobClass = &executionv1alpha1.ResolvedJobClassReference{ + Name: class.Name, UID: class.UID, ControllerName: class.Spec.ControllerName, + ParametersRef: executionv1alpha1.ParametersReference{Group: class.Spec.ParametersRef.Group, Kind: class.Spec.ParametersRef.Kind, Name: parameters.Name, UID: parameters.UID}, + } + job.Status.EffectiveResources = resources + job.Status.Execution = &executionv1alpha1.ExecutionStatus{ + Adapter: "kubernetes", + References: []executionv1alpha1.ExecutionReference{{Type: kubernetesadapter.ReferenceType, ID: string(backend.UID)}}, + } + meta.SetStatusCondition(&job.Status.Conditions, condition(job, executionv1alpha1.JobConditionAccepted, metav1.ConditionTrue, "Accepted", "Job accepted")) + meta.SetStatusCondition(&job.Status.Conditions, condition(job, executionv1alpha1.JobConditionScheduled, metav1.ConditionTrue, "BackendCreated", "Kubernetes Job created")) + meta.SetStatusCondition(&job.Status.Conditions, condition(job, executionv1alpha1.JobConditionSucceeded, metav1.ConditionUnknown, "Pending", "Waiting for task to start")) + job.Status.ObservedGeneration = job.Generation + return r.Status().Update(ctx, job) +} + +func (r *JobReconciler) observe(ctx context.Context, job *executionv1alpha1.Job, backend *batchv1.Job) error { + if job.Status.StartTime == nil && backend.Status.StartTime != nil { + job.Status.StartTime = backend.Status.StartTime.DeepCopy() + } + for _, backendCondition := range backend.Status.Conditions { + switch { + case backendCondition.Type == batchv1.JobComplete && backendCondition.Status == corev1.ConditionTrue: + completion := backend.Status.CompletionTime + if completion == nil { + now := metav1.NewTime(r.now()) + completion = &now + } + job.Status.CompletionTime = completion.DeepCopy() + job.Status.Result = &executionv1alpha1.JobResult{Reason: "Completed"} + return r.setCondition(ctx, job, metav1.Condition{Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionTrue, Reason: "Completed", Message: backendCondition.Message}) + case backendCondition.Type == batchv1.JobFailed && backendCondition.Status == corev1.ConditionTrue: + completion := metav1.NewTime(r.now()) + job.Status.CompletionTime = &completion + job.Status.Result = &executionv1alpha1.JobResult{Reason: "ProcessFailed"} + return r.setCondition(ctx, job, metav1.Condition{Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionFalse, Reason: "ProcessFailed", Message: backendCondition.Message}) + } + } + reason := "Pending" + message := "Waiting for task to start" + if backend.Status.StartTime != nil || backend.Status.Active > 0 { + reason = "Running" + message = "Task is running" + } + return r.setCondition(ctx, job, metav1.Condition{Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionUnknown, Reason: reason, Message: message}) +} + +func (r *JobReconciler) cancel(ctx context.Context, job *executionv1alpha1.Job) error { + backend := &batchv1.Job{} + key := types.NamespacedName{Namespace: job.Namespace, Name: job.Name} + err := r.Get(ctx, key, backend) + if err == nil { + if err := kubernetesadapter.ValidateOwnership(job, backend); err != nil { + return err + } + for _, backendCondition := range backend.Status.Conditions { + if (backendCondition.Type == batchv1.JobComplete || backendCondition.Type == batchv1.JobFailed) && + backendCondition.Status == corev1.ConditionTrue { + return r.observe(ctx, job, backend) + } + } + if err := r.Delete(ctx, backend, client.PropagationPolicy(metav1.DeletePropagationBackground)); err != nil && !apierrors.IsNotFound(err) { + return err + } + return nil + } + if !apierrors.IsNotFound(err) { + return err + } + now := metav1.NewTime(r.now()) + job.Status.CompletionTime = &now + job.Status.Result = &executionv1alpha1.JobResult{Reason: "Cancelled"} + return r.setCondition(ctx, job, metav1.Condition{Type: executionv1alpha1.JobConditionSucceeded, Status: metav1.ConditionFalse, Reason: "Cancelled", Message: "Execution cancelled"}) +} + +func (r *JobReconciler) finalize(ctx context.Context, job *executionv1alpha1.Job) error { + if !containsString(job.Finalizers, jobFinalizer) { + return nil + } + backend := &batchv1.Job{} + key := types.NamespacedName{Namespace: job.Namespace, Name: job.Name} + if err := r.Get(ctx, key, backend); err == nil { + if err := kubernetesadapter.ValidateOwnership(job, backend); err != nil { + return err + } + if err := r.Delete(ctx, backend, client.PropagationPolicy(metav1.DeletePropagationBackground)); err != nil && !apierrors.IsNotFound(err) { + return err + } + return nil + } else if !apierrors.IsNotFound(err) { + return err + } + job.Finalizers = removeString(job.Finalizers, jobFinalizer) + return r.Update(ctx, job) +} + +func (r *JobReconciler) setCondition(ctx context.Context, job *executionv1alpha1.Job, next metav1.Condition) error { + meta.SetStatusCondition(&job.Status.Conditions, condition(job, next.Type, next.Status, next.Reason, next.Message)) + job.Status.ObservedGeneration = job.Generation + return r.Status().Update(ctx, job) +} + +func condition(job *executionv1alpha1.Job, conditionType string, status metav1.ConditionStatus, reason, message string) metav1.Condition { + return metav1.Condition{Type: conditionType, Status: status, Reason: reason, Message: message, ObservedGeneration: job.Generation} +} + +func conditionTrue(conditions []metav1.Condition, conditionType string) bool { + current := meta.FindStatusCondition(conditions, conditionType) + return current != nil && current.Status == metav1.ConditionTrue +} + +func isTerminal(job *executionv1alpha1.Job) bool { + current := meta.FindStatusCondition(job.Status.Conditions, executionv1alpha1.JobConditionSucceeded) + return current != nil && (current.Status == metav1.ConditionTrue || current.Status == metav1.ConditionFalse) +} + +func applyResourceDefaults(requested, defaults executionv1alpha1.ExecutionResourceRequirements) executionv1alpha1.ExecutionResourceRequirements { + result := requested.DeepCopy() + if result.Requests.CPU == nil && defaults.Requests.CPU != nil { + result.Requests.CPU = copyQuantity(defaults.Requests.CPU) + } + if result.Requests.Memory == nil && defaults.Requests.Memory != nil { + result.Requests.Memory = copyQuantity(defaults.Requests.Memory) + } + if result.Limits.CPU == nil && defaults.Limits.CPU != nil { + result.Limits.CPU = copyQuantity(defaults.Limits.CPU) + } + if result.Limits.Memory == nil && defaults.Limits.Memory != nil { + result.Limits.Memory = copyQuantity(defaults.Limits.Memory) + } + return *result +} + +func copyQuantity(value *resource.Quantity) *resource.Quantity { + copy := value.DeepCopy() + return © +} + +func validateResources(resources executionv1alpha1.ExecutionResourceRequirements) error { + if resources.Requests.CPU != nil && resources.Limits.CPU != nil && resources.Requests.CPU.Cmp(*resources.Limits.CPU) > 0 { + return fmt.Errorf("CPU request must not exceed limit") + } + if resources.Requests.Memory != nil && resources.Limits.Memory != nil && resources.Requests.Memory.Cmp(*resources.Limits.Memory) > 0 { + return fmt.Errorf("memory request must not exceed limit") + } + return nil +} + +func containsString(values []string, target string) bool { + return slices.Contains(values, target) +} + +func removeString(values []string, target string) []string { + result := values[:0] + for _, value := range values { + if value != target { + result = append(result, value) + } + } + return result +} + +func (r *JobReconciler) now() time.Time { + if r.Now != nil { + return r.Now() + } + return time.Now() +} + +func (r *JobReconciler) SetupWithManager(manager ctrl.Manager) error { + return ctrl.NewControllerManagedBy(manager). + For(&executionv1alpha1.Job{}). + Owns(&batchv1.Job{}). + Named("execution-job"). + Complete(r) +} diff --git a/internal/controller/job_controller_test.go b/internal/controller/job_controller_test.go new file mode 100644 index 0000000..124ab94 --- /dev/null +++ b/internal/controller/job_controller_test.go @@ -0,0 +1,239 @@ +package controller + +import ( + "context" + "testing" + "time" + + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + kubernetesadapter "git.ddupan.top/panxiao81/ayatori/internal/adapter/kubernetes" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +const defaultClassName = "default" + +//nolint:modernize // controller-runtime and Kubernetes API structs expose promoted embedded fields. +func TestJobReconcilerKubernetesLifecycle(t *testing.T) { + ctx := context.Background() + now := time.Unix(1_700_000_000, 0) + reconciler, kubeClient := testReconciler(t, now, validObjects()...) + request := ctrl.Request{} + request.NamespacedName = types.NamespacedName{Namespace: "ci", Name: "hello"} + + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatalf("add finalizer: %v", err) + } + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatalf("create backend: %v", err) + } + + backend := &batchv1.Job{} + if err := kubeClient.Get(ctx, request.NamespacedName, backend); err != nil { + t.Fatalf("backend Job was not created: %v", err) + } + if backend.Labels[kubernetesadapter.JobUIDLabel] != "ayatori-job-uid" { + t.Fatalf("backend UID label = %q", backend.Labels[kubernetesadapter.JobUIDLabel]) + } + + job := getJob(t, ctx, kubeClient, request.NamespacedName) + if !conditionIs(job, executionv1alpha1.JobConditionAccepted, metav1.ConditionTrue) || + !conditionIs(job, executionv1alpha1.JobConditionScheduled, metav1.ConditionTrue) { + t.Fatalf("Job was not accepted and scheduled: %#v", job.Status.Conditions) + } + + started := metav1.NewTime(now.Add(time.Minute)) + backend.Status.StartTime = &started + backend.Status.Active = 1 + if err := kubeClient.Status().Update(ctx, backend); err != nil { + t.Fatalf("set backend running: %v", err) + } + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatalf("observe running backend: %v", err) + } + job = getJob(t, ctx, kubeClient, request.NamespacedName) + if job.Status.StartTime == nil || !conditionIs(job, executionv1alpha1.JobConditionSucceeded, metav1.ConditionUnknown) { + t.Fatalf("running state was not observed: %#v", job.Status) + } + + completed := metav1.NewTime(now.Add(2 * time.Minute)) + backend = &batchv1.Job{} + if err := kubeClient.Get(ctx, request.NamespacedName, backend); err != nil { + t.Fatal(err) + } + backend.Status.Active = 0 + backend.Status.CompletionTime = &completed + backend.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue, Reason: "Completed"}} + if err := kubeClient.Status().Update(ctx, backend); err != nil { + t.Fatalf("set backend complete: %v", err) + } + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatalf("observe completed backend: %v", err) + } + job = getJob(t, ctx, kubeClient, request.NamespacedName) + if !conditionIs(job, executionv1alpha1.JobConditionSucceeded, metav1.ConditionTrue) || job.Status.CompletionTime == nil { + t.Fatalf("terminal state was not observed: %#v", job.Status) + } +} + +//nolint:modernize // controller-runtime Request exposes NamespacedName as a promoted embedded field. +func TestJobReconcilerRejectsMissingClass(t *testing.T) { + ctx := context.Background() + job := validObjects()[3].(*executionv1alpha1.Job).DeepCopy() + job.Spec.JobClassName = "missing" + reconciler, kubeClient := testReconciler(t, time.Now(), validObjects()[0], validObjects()[1], job) + request := ctrl.Request{} + request.NamespacedName = types.NamespacedName{Namespace: job.Namespace, Name: job.Name} + + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatal(err) + } + result, err := reconciler.Reconcile(ctx, request) + if err != nil { + t.Fatal(err) + } + if result.RequeueAfter == 0 { + t.Fatal("missing JobClass did not schedule a retry") + } + stored := getJob(t, ctx, kubeClient, request.NamespacedName) + accepted := meta.FindStatusCondition(stored.Status.Conditions, executionv1alpha1.JobConditionAccepted) + if accepted == nil || accepted.Status != metav1.ConditionFalse || accepted.Reason != "JobClassNotFound" { + t.Fatalf("unexpected Accepted condition: %#v", accepted) + } +} + +//nolint:modernize // controller-runtime Request exposes NamespacedName as a promoted embedded field. +func TestJobReconcilerObservesExistingExecutionWithoutJobClass(t *testing.T) { + ctx := context.Background() + now := time.Unix(1_700_000_000, 0) + job := validObjects()[3].(*executionv1alpha1.Job).DeepCopy() + job.Finalizers = []string{jobFinalizer} + job.Status.Execution = &executionv1alpha1.ExecutionStatus{Adapter: "kubernetes"} + backend := kubernetesadapter.BuildJob(job, validObjects()[4].(*executionv1alpha1.KubernetesExecutionParameters), executionv1alpha1.ExecutionResourceRequirements{}) + backend.Status.StartTime = &metav1.Time{Time: now} + backend.Status.Active = 1 + reconciler, kubeClient := testReconciler(t, now, job, backend) + request := ctrl.Request{NamespacedName: types.NamespacedName{Namespace: job.Namespace, Name: job.Name}} + + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatal(err) + } + stored := getJob(t, ctx, kubeClient, request.NamespacedName) + if stored.Status.StartTime == nil || !conditionIs(stored, executionv1alpha1.JobConditionSucceeded, metav1.ConditionUnknown) { + t.Fatalf("existing execution was not observed without its JobClass: %#v", stored.Status) + } +} + +//nolint:modernize // controller-runtime Request exposes NamespacedName as a promoted embedded field. +func TestJobReconcilerCancelsBeforeScheduling(t *testing.T) { + ctx := context.Background() + job := validObjects()[3].(*executionv1alpha1.Job).DeepCopy() + job.Spec.DesiredState = executionv1alpha1.JobDesiredStateCancelled + reconciler, kubeClient := testReconciler(t, time.Unix(1_700_000_000, 0), job) + request := ctrl.Request{NamespacedName: types.NamespacedName{Namespace: job.Namespace, Name: job.Name}} + + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatal(err) + } + stored := getJob(t, ctx, kubeClient, request.NamespacedName) + condition := meta.FindStatusCondition(stored.Status.Conditions, executionv1alpha1.JobConditionSucceeded) + if condition == nil || condition.Status != metav1.ConditionFalse || condition.Reason != "Cancelled" { + t.Fatalf("unexpected cancellation condition: %#v", condition) + } + if stored.Status.CompletionTime == nil { + t.Fatal("cancelled Job has no completionTime") + } +} + +//nolint:modernize // controller-runtime Request exposes NamespacedName as a promoted embedded field. +func TestJobReconcilerKeepsConfirmedSuccessDuringCancellation(t *testing.T) { + ctx := context.Background() + now := time.Unix(1_700_000_000, 0) + job := validObjects()[3].(*executionv1alpha1.Job).DeepCopy() + job.Spec.DesiredState = executionv1alpha1.JobDesiredStateCancelled + job.Finalizers = []string{jobFinalizer} + backend := kubernetesadapter.BuildJob(job, validObjects()[4].(*executionv1alpha1.KubernetesExecutionParameters), executionv1alpha1.ExecutionResourceRequirements{}) + backend.Status.CompletionTime = &metav1.Time{Time: now} + backend.Status.Conditions = []batchv1.JobCondition{{Type: batchv1.JobComplete, Status: corev1.ConditionTrue}} + reconciler, kubeClient := testReconciler(t, now, job, backend) + request := ctrl.Request{NamespacedName: types.NamespacedName{Namespace: job.Namespace, Name: job.Name}} + + if _, err := reconciler.Reconcile(ctx, request); err != nil { + t.Fatal(err) + } + stored := getJob(t, ctx, kubeClient, request.NamespacedName) + if !conditionIs(stored, executionv1alpha1.JobConditionSucceeded, metav1.ConditionTrue) { + t.Fatalf("confirmed success was overwritten by cancellation: %#v", stored.Status.Conditions) + } + if err := kubeClient.Get(ctx, request.NamespacedName, &batchv1.Job{}); err != nil { + t.Fatalf("successful backend was deleted: %v", err) + } +} + +//nolint:modernize // Kubernetes API structs expose ObjectMeta through embedded TypeMeta fields. +func validObjects() []client.Object { + return []client.Object{ + &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "ci", Labels: map[string]string{"execution": "enabled"}}}, + &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: "runner", Namespace: "ci"}}, + &executionv1alpha1.JobClass{ + ObjectMeta: metav1.ObjectMeta{Name: defaultClassName, UID: types.UID("class-uid")}, + Spec: executionv1alpha1.JobClassSpec{ + ControllerName: kubernetesadapter.ControllerName, + ParametersRef: executionv1alpha1.ParametersReference{Group: executionv1alpha1.GroupVersion.Group, Kind: "KubernetesExecutionParameters", Name: defaultClassName}, + AllowedNamespaces: &metav1.LabelSelector{MatchLabels: map[string]string{"execution": "enabled"}}, + }, + }, + &executionv1alpha1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "hello", Namespace: "ci", UID: types.UID("ayatori-job-uid")}, + Spec: executionv1alpha1.JobSpec{ + JobClassName: defaultClassName, DesiredState: executionv1alpha1.JobDesiredStateRunning, + Task: executionv1alpha1.TaskSpec{Image: "alpine:3.22", Command: []string{"true"}}, + }, + }, + &executionv1alpha1.KubernetesExecutionParameters{ + ObjectMeta: metav1.ObjectMeta{Name: defaultClassName, UID: types.UID("parameters-uid")}, + Spec: executionv1alpha1.KubernetesExecutionParametersSpec{ServiceAccountName: "runner", ImagePullPolicy: corev1.PullIfNotPresent}, + }, + } +} + +func testReconciler(t *testing.T, now time.Time, objects ...client.Object) (*JobReconciler, client.Client) { + t.Helper() + scheme := runtime.NewScheme() + if err := corev1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + if err := batchv1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + if err := executionv1alpha1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + kubeClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithStatusSubresource(&executionv1alpha1.Job{}, &batchv1.Job{}). + WithObjects(objects...). + Build() + return &JobReconciler{Client: kubeClient, Now: func() time.Time { return now }}, kubeClient +} + +func getJob(t *testing.T, ctx context.Context, kubeClient client.Client, key types.NamespacedName) *executionv1alpha1.Job { + t.Helper() + job := &executionv1alpha1.Job{} + if err := kubeClient.Get(ctx, key, job); err != nil { + t.Fatal(err) + } + return job +} + +func conditionIs(job *executionv1alpha1.Job, conditionType string, status metav1.ConditionStatus) bool { + condition := meta.FindStatusCondition(job.Status.Conditions, conditionType) + return condition != nil && condition.Status == status +} -- 2.54.0 From 06bc54e3cf1eed0860b32dbd911d9ca4fc0f840a Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Fri, 18 Sep 2026 18:30:14 +0000 Subject: [PATCH 2/3] test: exercise controller with envtest --- .../job_controller_integration_test.go | 202 ++++++++++++++++++ internal/controller/job_controller_test.go | 20 +- 2 files changed, 215 insertions(+), 7 deletions(-) create mode 100644 internal/controller/job_controller_integration_test.go diff --git a/internal/controller/job_controller_integration_test.go b/internal/controller/job_controller_integration_test.go new file mode 100644 index 0000000..8606a00 --- /dev/null +++ b/internal/controller/job_controller_integration_test.go @@ -0,0 +1,202 @@ +package controller + +import ( + "context" + "fmt" + "os" + "path/filepath" + "testing" + "time" + + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" + kubernetesadapter "git.ddupan.top/panxiao81/ayatori/internal/adapter/kubernetes" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/envtest" + metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" +) + +const ( + integrationNamespace = "controller-integration" + integrationClass = "integration" +) + +//nolint:modernize // Kubernetes API structs expose ObjectMeta through embedded TypeMeta fields. +func TestJobControllerIntegration(t *testing.T) { + if os.Getenv("KUBEBUILDER_ASSETS") == "" { + t.Skip("KUBEBUILDER_ASSETS is unset; run make test to execute controller integration tests") + } + + scheme := runtime.NewScheme() + for _, addToScheme := range []func(*runtime.Scheme) error{ + corev1.AddToScheme, + batchv1.AddToScheme, + executionv1alpha1.AddToScheme, + } { + if err := addToScheme(scheme); err != nil { + t.Fatal(err) + } + } + + crdPath, err := filepath.Abs("../../config/crd/bases") + if err != nil { + t.Fatal(err) + } + environment := &envtest.Environment{CRDDirectoryPaths: []string{crdPath}} + config, err := environment.Start() + if err != nil { + t.Fatalf("start envtest: %v", err) + } + t.Cleanup(func() { + if err := environment.Stop(); err != nil { + t.Errorf("stop envtest: %v", err) + } + }) + + manager, err := ctrl.NewManager(config, ctrl.Options{ + Scheme: scheme, + Metrics: metricsserver.Options{BindAddress: "0"}, + }) + if err != nil { + t.Fatal(err) + } + if err := (&JobReconciler{Client: manager.GetClient()}).SetupWithManager(manager); err != nil { + t.Fatal(err) + } + + managerContext, cancelManager := context.WithCancel(context.Background()) + t.Cleanup(cancelManager) + managerErrors := make(chan error, 1) + go func() { + managerErrors <- manager.Start(managerContext) + }() + if !manager.GetCache().WaitForCacheSync(managerContext) { + t.Fatal("manager cache did not synchronize") + } + + directClient, err := client.New(config, client.Options{Scheme: scheme}) + if err != nil { + t.Fatal(err) + } + ctx := context.Background() + objects := []client.Object{ + &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: integrationNamespace, Labels: map[string]string{testLabelKey: testLabelEnabled}}}, + &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: testSAName, Namespace: integrationNamespace}}, + &executionv1alpha1.KubernetesExecutionParameters{ + ObjectMeta: metav1.ObjectMeta{Name: integrationClass}, + Spec: executionv1alpha1.KubernetesExecutionParametersSpec{ + ServiceAccountName: testSAName, + ImagePullPolicy: corev1.PullIfNotPresent, + }, + }, + &executionv1alpha1.JobClass{ + ObjectMeta: metav1.ObjectMeta{Name: integrationClass}, + Spec: executionv1alpha1.JobClassSpec{ + ControllerName: kubernetesadapter.ControllerName, + ParametersRef: executionv1alpha1.ParametersReference{ + Group: executionv1alpha1.GroupVersion.Group, + Kind: "KubernetesExecutionParameters", + Name: integrationClass, + }, + AllowedNamespaces: &metav1.LabelSelector{MatchLabels: map[string]string{testLabelKey: testLabelEnabled}}, + }, + }, + } + for _, object := range objects { + if err := directClient.Create(ctx, object); err != nil { + t.Fatalf("create %T: %v", object, err) + } + } + + job := &executionv1alpha1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: testJobName, Namespace: integrationNamespace}, + Spec: executionv1alpha1.JobSpec{ + JobClassName: integrationClass, + DesiredState: executionv1alpha1.JobDesiredStateRunning, + Task: executionv1alpha1.TaskSpec{Image: "alpine:3.22", Command: []string{"true"}}, + }, + } + if err := directClient.Create(ctx, job); err != nil { + t.Fatal(err) + } + + backend := &batchv1.Job{} + eventually(t, 10*time.Second, func() (bool, error) { + err := directClient.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: job.Name}, backend) + return err == nil, client.IgnoreNotFound(err) + }) + if backend.Labels[kubernetesadapter.JobUIDLabel] != string(job.UID) { + t.Fatalf("backend identity label = %q, want %q", backend.Labels[kubernetesadapter.JobUIDLabel], job.UID) + } + + eventually(t, 10*time.Second, func() (bool, error) { + if err := directClient.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: job.Name}, job); err != nil { + return false, err + } + return conditionStatus(job, executionv1alpha1.JobConditionScheduled) == metav1.ConditionTrue, nil + }) + + completed := metav1.Now() + if err := directClient.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: job.Name}, backend); err != nil { + t.Fatal(err) + } + backend.Status.StartTime = &completed + backend.Status.CompletionTime = &completed + backend.Status.Conditions = []batchv1.JobCondition{ + {Type: batchv1.JobSuccessCriteriaMet, Status: corev1.ConditionTrue, Reason: "CompletionsReached"}, + {Type: batchv1.JobComplete, Status: corev1.ConditionTrue, Reason: "Completed"}, + } + if err := directClient.Status().Update(ctx, backend); err != nil { + t.Fatal(err) + } + + eventually(t, 10*time.Second, func() (bool, error) { + if err := directClient.Get(ctx, types.NamespacedName{Namespace: job.Namespace, Name: job.Name}, job); err != nil { + return false, err + } + return conditionStatus(job, executionv1alpha1.JobConditionSucceeded) == metav1.ConditionTrue, nil + }) + if job.Status.StartTime == nil || job.Status.CompletionTime == nil || job.Status.Execution == nil { + t.Fatalf("controller did not persist execution status: %#v", job.Status) + } + + cancelManager() + select { + case err := <-managerErrors: + if err != nil { + t.Fatalf("manager stopped with error: %v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("manager did not stop") + } +} + +func eventually(t *testing.T, timeout time.Duration, check func() (bool, error)) { + t.Helper() + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + ready, err := check() + if err != nil { + t.Fatal(err) + } + if ready { + return + } + time.Sleep(100 * time.Millisecond) + } + t.Fatal(fmt.Errorf("condition was not met within %s", timeout)) +} + +func conditionStatus(job *executionv1alpha1.Job, conditionType string) metav1.ConditionStatus { + condition := meta.FindStatusCondition(job.Status.Conditions, conditionType) + if condition == nil { + return metav1.ConditionUnknown + } + return condition.Status +} diff --git a/internal/controller/job_controller_test.go b/internal/controller/job_controller_test.go index 124ab94..fca0afc 100644 --- a/internal/controller/job_controller_test.go +++ b/internal/controller/job_controller_test.go @@ -18,7 +18,13 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client/fake" ) -const defaultClassName = "default" +const ( + defaultClassName = "default" + testJobName = "hello" + testSAName = "runner" + testLabelKey = "execution" + testLabelEnabled = "enabled" +) //nolint:modernize // controller-runtime and Kubernetes API structs expose promoted embedded fields. func TestJobReconcilerKubernetesLifecycle(t *testing.T) { @@ -26,7 +32,7 @@ func TestJobReconcilerKubernetesLifecycle(t *testing.T) { now := time.Unix(1_700_000_000, 0) reconciler, kubeClient := testReconciler(t, now, validObjects()...) request := ctrl.Request{} - request.NamespacedName = types.NamespacedName{Namespace: "ci", Name: "hello"} + request.NamespacedName = types.NamespacedName{Namespace: "ci", Name: testJobName} if _, err := reconciler.Reconcile(ctx, request); err != nil { t.Fatalf("add finalizer: %v", err) @@ -180,18 +186,18 @@ func TestJobReconcilerKeepsConfirmedSuccessDuringCancellation(t *testing.T) { //nolint:modernize // Kubernetes API structs expose ObjectMeta through embedded TypeMeta fields. func validObjects() []client.Object { return []client.Object{ - &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "ci", Labels: map[string]string{"execution": "enabled"}}}, - &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: "runner", Namespace: "ci"}}, + &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: "ci", Labels: map[string]string{testLabelKey: testLabelEnabled}}}, + &corev1.ServiceAccount{ObjectMeta: metav1.ObjectMeta{Name: testSAName, Namespace: "ci"}}, &executionv1alpha1.JobClass{ ObjectMeta: metav1.ObjectMeta{Name: defaultClassName, UID: types.UID("class-uid")}, Spec: executionv1alpha1.JobClassSpec{ ControllerName: kubernetesadapter.ControllerName, ParametersRef: executionv1alpha1.ParametersReference{Group: executionv1alpha1.GroupVersion.Group, Kind: "KubernetesExecutionParameters", Name: defaultClassName}, - AllowedNamespaces: &metav1.LabelSelector{MatchLabels: map[string]string{"execution": "enabled"}}, + AllowedNamespaces: &metav1.LabelSelector{MatchLabels: map[string]string{testLabelKey: testLabelEnabled}}, }, }, &executionv1alpha1.Job{ - ObjectMeta: metav1.ObjectMeta{Name: "hello", Namespace: "ci", UID: types.UID("ayatori-job-uid")}, + ObjectMeta: metav1.ObjectMeta{Name: testJobName, Namespace: "ci", UID: types.UID("ayatori-job-uid")}, Spec: executionv1alpha1.JobSpec{ JobClassName: defaultClassName, DesiredState: executionv1alpha1.JobDesiredStateRunning, Task: executionv1alpha1.TaskSpec{Image: "alpine:3.22", Command: []string{"true"}}, @@ -199,7 +205,7 @@ func validObjects() []client.Object { }, &executionv1alpha1.KubernetesExecutionParameters{ ObjectMeta: metav1.ObjectMeta{Name: defaultClassName, UID: types.UID("parameters-uid")}, - Spec: executionv1alpha1.KubernetesExecutionParametersSpec{ServiceAccountName: "runner", ImagePullPolicy: corev1.PullIfNotPresent}, + Spec: executionv1alpha1.KubernetesExecutionParametersSpec{ServiceAccountName: testSAName, ImagePullPolicy: corev1.PullIfNotPresent}, }, } } -- 2.54.0 From f76978ca17e01860a172e6e660636d44583c92eb Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Fri, 18 Sep 2026 18:33:11 +0000 Subject: [PATCH 3/3] docs: require integration tests for controllers --- AGENTS.md | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/AGENTS.md b/AGENTS.md index c40b995..51e2f16 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -8,6 +8,16 @@ 理由。 - 不要引入统一包装所有能力的 Application CRD;应用应直接组合正交的平台资源。 - 所有 controller 必须考虑幂等、observe、finalizer、conditions、删除策略和恢复行为。 +- Ayatori 会联动 Kubernetes API、虚拟化、存储、网络及其他外部控制面;集成测试是功能完成 + 标准的一部分,不得仅凭 fake client 或 mock 测试宣告 controller、adapter 或生命周期变更完成。 +- 测试应按风险分层:纯领域规则使用快速单元测试;API schema、CEL、status subresource、 + watch/cache、owner reference 和 reconcile 事件链使用 envtest;需要 scheduler、kubelet、网络、 + 存储或真实后端行为的路径在 Dev 集群或对应后端环境执行端到端测试。 +- fake client 适合穷举状态机和错误分支,但它不会完整执行 API server defaulting、validation、 + resourceVersion、garbage collection 或新版 Kubernetes 约束;涉及这些语义时必须增加真实 API + server 测试。跨 adapter 的共同契约应使用同一套 contract tests,避免各实现产生语义漂移。 +- 集成测试必须覆盖正常路径以及幂等重试、controller 重启、依赖稍后出现、删除/finalizer、 + 后端结果不确定和并发竞态等恢复路径;无法在当前层测试的部分要明确记录由哪一层验证。 - Secret、token、kubeconfig 及具体生产凭据不得提交到仓库。 - `deploy/dev/` 与 `deploy/prod/` 使用相同制品;生产版本只通过 promotion 更新。 - 内部专用不构成降低测试、版本、恢复、安全和可审计要求的理由。 -- 2.54.0