// Изменено: 2026-03-20 (function-service-split: FunctionReconciler — только build pipeline) // FunctionReconciler — контроллер Function CRD (sless_function = oneshot/Job). // Функция = код который выполняется ОДИН РАЗ через k8s Job при каждом вызове. // Нет Deployment, нет постоянного URL. Вызов — через FunctionJob или invoke API. // Reconciler отвечает только за: // 1. Сборку Docker-образа через kaniko (Pending → Building → Ready/Failed) // 2. Очистку ресурсов при удалении (kaniko Job) // Deployment/Service/Ingress — в ServiceReconciler (sless_service). package controllers import ( "bytes" "context" "fmt" "io" "strings" "time" corev1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/client-go/kubernetes" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/log" slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1" "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/builder" "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/internal/harbor" ) // FunctionReconciler reconciles a Function object type FunctionReconciler struct { client.Client Scheme *runtime.Scheme Builder *builder.Builder KubeClient kubernetes.Interface // typed client для чтения логов build-подов RegistrySecret string // имя Secret с docker credentials (для imagePullSecrets в подах функций) OperatorNamespace string // namespace оператора — откуда копируем RegistrySecret в sless-fn-* HarborClient *harbor.Client // nil — Harbor не используется, EnsureProject пропускается } //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions/status,verbs=get;update;patch //+kubebuilder:rbac:groups=sless.kube5s.ru,resources=functions/finalizers,verbs=update //+kubebuilder:rbac:groups=apps,resources=deployments,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups=batch,resources=jobs,verbs=get;list;watch;create;update;patch;delete //+kubebuilder:rbac:groups="",resources=namespaces,verbs=get;list;watch;create //+kubebuilder:rbac:groups="",resources=secrets,verbs=get;create //+kubebuilder:rbac:groups="",resources=events,verbs=create;patch // Reconcile — главный цикл управления Function. // Логика: читаем текущее состояние → определяем что нужно сделать → делаем. func (r *FunctionReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { logger := log.FromContext(ctx) // Читаем Function объект из k8s fn := &slessv1alpha1.Function{} if err := r.Get(ctx, req.NamespacedName, fn); err != nil { if errors.IsNotFound(err) { // Объект удалён — ничего делать не надо, finalizer уже отработал return ctrl.Result{}, nil } return ctrl.Result{}, fmt.Errorf("get function: %w", err) } // Обрабатываем удаление через finalizer if !fn.DeletionTimestamp.IsZero() { return r.handleDeletion(ctx, fn) } // Добавляем finalizer при первом создании чтобы обработать удаление if !containsString(fn.Finalizers, finalizerName) { fn.Finalizers = append(fn.Finalizers, finalizerName) if err := r.Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("add finalizer: %w", err) } return ctrl.Result{Requeue: true}, nil } // Есть новый код (s3Key изменился по сравнению с последним запуском сборки) // и сборка сейчас не идёт — запускаем. Это единственное место где решается "нужна ли сборка". // Почему аннотация а не phase: phase может быть Failed/Ready от прошлого кода; // новый upload меняет Spec.S3Key → контроллер сам понимает что нужно пересобрать. builtKey := fn.Annotations["sless.kube5s.ru/last-built-s3key"] needsBuild := fn.Spec.S3Key != "" && builtKey != fn.Spec.S3Key if needsBuild && fn.Status.Phase != slessv1alpha1.FunctionPhaseBuilding { logger.Info("starting build", "function", fn.Name) return r.startBuild(ctx, fn) } switch fn.Status.Phase { case slessv1alpha1.FunctionPhaseBuilding: return r.checkBuild(ctx, fn) case slessv1alpha1.FunctionPhaseReady: // Function = oneshot. После успешной сборки образ готов — Deployment не создаём. // Вызов через FunctionJob или invoke API (Job per call). return ctrl.Result{}, nil } return ctrl.Result{}, nil } const finalizerName = "sless.kube5s.ru/finalizer" // startBuild проверяет наличие образа в registry и либо пропускает сборку, // либо запускает kaniko Job. Идемпотентность: если код не менялся (тег = hash s3Key), // образ уже в registry → deploy без пересборки. // Критически важно: СНАЧАЛА сохраняем last-built-s3key аннотацию, ПОТОМ status. func (r *FunctionReconciler) startBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { imageRef := r.Builder.ImageRef(fn.Namespace, fn.Name, fn.Spec.S3Key) // Проверяем: образ с этим тегом уже существует в registry? // Если да — пропускаем kaniko, сразу переходим в Ready. // Если registry недоступен — requeue, не запускаем сборку (kaniko тоже упадёт). exists, err := r.Builder.ImageExists(ctx, imageRef) if err != nil { return ctrl.Result{RequeueAfter: 10 * time.Second}, fmt.Errorf("check image exists: %w", err) } if exists { logger := log.FromContext(ctx) logger.Info("image already exists in registry, skipping build", "imageRef", imageRef) if fn.Annotations == nil { fn.Annotations = map[string]string{} } fn.Annotations["sless.kube5s.ru/last-built-s3key"] = fn.Spec.S3Key if err := r.Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update annotations (cache hit): %w", err) } fn.Status.Phase = slessv1alpha1.FunctionPhaseReady fn.Status.ImageRef = imageRef fn.Status.Message = "Image restored from registry cache" if err := r.Status().Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update status (cache hit): %w", err) } return ctrl.Result{}, nil } jobName, err := r.Builder.Build(ctx, fn.Namespace, fn.Name, fn.Spec.S3Key) if err != nil { return r.setFailed(ctx, fn, fmt.Sprintf("failed to start build: %v", err)) } // Обновляем аннотации ПЕРВЫМИ — это idempotency guard. // Как только last-built-s3key == spec.S3Key, дальнейшие reconcile не будут // вызывать startBuild снова, даже если status ещё не обновился. if fn.Annotations == nil { fn.Annotations = map[string]string{} } fn.Annotations["sless.kube5s.ru/build-job"] = jobName fn.Annotations["sless.kube5s.ru/last-built-s3key"] = fn.Spec.S3Key if err := r.Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update build annotations: %w", err) } // Обновляем статус после аннотаций fn.Status.Phase = slessv1alpha1.FunctionPhaseBuilding fn.Status.Message = "Building image: " + jobName if err := r.Status().Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update status to building: %w", err) } return ctrl.Result{RequeueAfter: 10 * time.Second}, nil } // checkBuild проверяет статус build Job'а. func (r *FunctionReconciler) checkBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { jobName := fn.Annotations["sless.kube5s.ru/build-job"] if jobName == "" { // Аннотация потерялась — перезапускаем сборку fn.Status.Phase = slessv1alpha1.FunctionPhasePending _ = r.Status().Update(ctx, fn) return ctrl.Result{Requeue: true}, nil } status, err := r.Builder.JobStatus(ctx, jobName) if err != nil { return ctrl.Result{}, fmt.Errorf("check build job: %w", err) } switch status { case "running": // Ещё идёт — проверим через 10 секунд return ctrl.Result{RequeueAfter: 10 * time.Second}, nil case "succeeded": imageRef := r.Builder.ImageRef(fn.Namespace, fn.Name, fn.Spec.S3Key) fn.Status.Phase = slessv1alpha1.FunctionPhaseReady fn.Status.ImageRef = imageRef fn.Status.Message = "" now := metav1.Now() fn.Status.LastBuiltAt = &now if err := r.Status().Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update status to ready: %w", err) } // Чистим завершённый Job _ = r.Builder.Cleanup(ctx, jobName) return ctrl.Result{Requeue: true}, nil case "failed": // Захватываем логи build-пода чтобы разработчик видел причину ошибки (pip error и т.д.). logs := getBuildPodLogs(ctx, r.KubeClient, r.OperatorNamespace, jobName) msg := "build job failed" if logs != "" { msg = "build job failed:\n" + logs } return r.setFailed(ctx, fn, msg) } return ctrl.Result{RequeueAfter: 10 * time.Second}, nil } // handleDeletion обрабатывает удаление Function: убивает kaniko Job и убирает finalizer. // Deployment/Service/Ingress Function не создаёт — они принадлежат Service CRD. func (r *FunctionReconciler) handleDeletion(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { // Убиваем kaniko Job если сборка шла в момент удаления if jobName := fn.Annotations["sless.kube5s.ru/build-job"]; jobName != "" { _ = r.Builder.Cleanup(ctx, jobName) } fn.Finalizers = removeString(fn.Finalizers, finalizerName) if err := r.Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("remove finalizer: %w", err) } return ctrl.Result{}, nil } // setFailed переводит функцию в фазу Failed с сообщением об ошибке. func (r *FunctionReconciler) setFailed(ctx context.Context, fn *slessv1alpha1.Function, msg string) (ctrl.Result, error) { fn.Status.Phase = slessv1alpha1.FunctionPhaseFailed fn.Status.Message = msg if err := r.Status().Update(ctx, fn); err != nil { return ctrl.Result{}, fmt.Errorf("update status to failed: %w", err) } return ctrl.Result{}, nil } func containsString(slice []string, s string) bool { for _, v := range slice { if v == s { return true } } return false } func removeString(slice []string, s string) []string { var result []string for _, v := range slice { if v != s { result = append(result, v) } } return result } // SetupWithManager sets up the controller with the Manager. func (r *FunctionReconciler) SetupWithManager(mgr ctrl.Manager) error { return ctrl.NewControllerManagedBy(mgr). For(&slessv1alpha1.Function{}). Complete(r) } // getBuildPodLogs возвращает логи (stderr+stdout) пода kaniko build Job'а. // Используется чтобы пробросить ошибку pip/kaniko в Function.Status.Message. // Возвращает не более 50 последних строк — достаточно для диагностики, не засоряет CRD. // Если логи недоступны — возвращает пустую строку (caller покажет generic msg). func getBuildPodLogs(ctx context.Context, kube kubernetes.Interface, namespace, jobName string) string { if kube == nil { return "" } pods, err := kube.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{ LabelSelector: "job-name=" + jobName, }) if err != nil || len(pods.Items) == 0 { return "" } req := kube.CoreV1().Pods(namespace).GetLogs(pods.Items[0].Name, &corev1.PodLogOptions{}) stream, err := req.Stream(ctx) if err != nil { return "" } defer stream.Close() buf := new(bytes.Buffer) _, _ = io.Copy(buf, stream) raw := strings.TrimSpace(buf.String()) if raw == "" { return "" } // Оставляем последние 50 строк — ошибки pip всегда в конце вывода. lines := strings.Split(raw, "\n") if len(lines) > 50 { lines = lines[len(lines)-50:] } return strings.Join(lines, "\n") }