// Изменено: 2026-03-11 // FunctionReconciler — основной контроллер оператора. // Следит за CRD Function и управляет lifecycle функции: // Pending → Building (запуск kaniko Job) → Ready (образ собран, Deployment создан) / Failed // Reconcile вызывается k8s при любом изменении Function объекта. package controllers import ( "bytes" "context" "fmt" "io" "sort" "strings" "time" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" netv1 "k8s.io/api/networking/v1" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" 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: return r.ensureDeployment(ctx, fn) } return ctrl.Result{}, nil } const finalizerName = "sless.kube5s.ru/finalizer" // startBuild запускает kaniko Job и помечает функцию как Building. // Критически важно: СНАЧАЛА сохраняем last-built-s3key аннотацию, ПОТОМ status. // Это предотвращает повторный запуск сборки при параллельных reconcile — // следующий reconcile увидит last-built-s3key == spec.S3Key и не войдёт в startBuild. func (r *FunctionReconciler) startBuild(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { 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 } // ensureDeployment создаёт или обновляет Deployment для HTTP функции. // Deployment запускается в отдельном namespace sless-fn-{namespace}. func (r *FunctionReconciler) ensureDeployment(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { deployNS := "sless-fn-" + fn.Namespace // Создаём namespace для функций если не существует ns := &corev1.Namespace{} if err := r.Get(ctx, client.ObjectKey{Name: deployNS}, ns); err != nil { if errors.IsNotFound(err) { ns = &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: deployNS}} if err := r.Create(ctx, ns); err != nil { return ctrl.Result{}, fmt.Errorf("create function namespace: %w", err) } // Создаём Harbor-проект для namespace сразу при создании k8s NS (best-effort). // Если не удалось — EnsureProject повторит вызов внутри Build(). if r.HarborClient != nil { if err := r.HarborClient.EnsureProject(ctx, fn.Namespace); err != nil { log.FromContext(ctx).Error(err, "harbor ensure project on ns create", "project", fn.Namespace) } } } else { return ctrl.Result{}, fmt.Errorf("get function namespace: %w", err) } } // Обеспечиваем наличие registry pull-секрета в namespace функций. // Без него kubelet не сможет pull-нуть private образ из Harbor. if r.RegistrySecret != "" && r.OperatorNamespace != "" { if err := r.ensureRegistrySecret(ctx, deployNS); err != nil { // Не фатальная ошибка — логируем, но продолжаем log.FromContext(ctx).Error(err, "failed to ensure registry secret", "ns", deployNS) } } desired := r.buildDeployment(fn, deployNS) existing := &appsv1.Deployment{} err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existing) if errors.IsNotFound(err) { if err := r.Create(ctx, desired); err != nil { return ctrl.Result{}, fmt.Errorf("create deployment: %w", err) } return ctrl.Result{}, nil } if err != nil { return ctrl.Result{}, fmt.Errorf("get deployment: %w", err) } // Обновляем образ, env и imagePullSecrets при пересборке или изменении конфига. // Тег образа уникален per build (sha256 от s3Key) → imagePullPolicy: IfNotPresent // корректно подтягивает новый образ без дополнительных хаков. // Env обновляем целиком — иначе изменение entrypoint/env_vars не применяется. existing.Spec.Template.Spec.Containers[0].Image = fn.Status.ImageRef existing.Spec.Template.Spec.Containers[0].Env = desired.Spec.Template.Spec.Containers[0].Env existing.Spec.Template.Spec.ImagePullSecrets = desired.Spec.Template.Spec.ImagePullSecrets if err := r.Update(ctx, existing); err != nil { return ctrl.Result{}, fmt.Errorf("update deployment: %w", err) } return ctrl.Result{}, nil } // Изменено: 2026-03-11// buildDeployment формирует Deployment манифест для функции. func (r *FunctionReconciler) buildDeployment(fn *slessv1alpha1.Function, namespace string) *appsv1.Deployment { replicas := int32(1) envVars := []corev1.EnvVar{ // SLESS_ENTRYPOINT сообщает server.py/server.js какой файл и функцию загружать. // Формат: "module-name.funcName" (например: handler-http.handle) {Name: "SLESS_ENTRYPOINT", Value: fn.Spec.Entrypoint}, } // Сортируем ключи env vars для стабильного порядка в Pod spec. // map range в Go — недетерминирован: разный порядок при каждом вызове. // Нестабильный порядок → k8s видит изменение контейнера → лишние rollout'ы. keys := make([]string, 0, len(fn.Spec.Env)) for k := range fn.Spec.Env { keys = append(keys, k) } sort.Strings(keys) for _, k := range keys { envVars = append(envVars, corev1.EnvVar{Name: k, Value: fn.Spec.Env[k]}) } return &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ Name: fn.Name, Namespace: namespace, Labels: map[string]string{"app": fn.Name, "managed-by": "sless"}, }, Spec: appsv1.DeploymentSpec{ Replicas: &replicas, Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": fn.Name}}, Template: corev1.PodTemplateSpec{ ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{"app": fn.Name}}, Spec: corev1.PodSpec{ Containers: []corev1.Container{ { Name: fn.Name, Image: fn.Status.ImageRef, Env: envVars, Resources: corev1.ResourceRequirements{ Limits: corev1.ResourceList{ corev1.ResourceMemory: resource.MustParse(fmt.Sprintf("%dMi", fn.Spec.MemoryMB)), }, }, }, }, ImagePullSecrets: func() []corev1.LocalObjectReference { if r.RegistrySecret != "" { return []corev1.LocalObjectReference{{Name: r.RegistrySecret}} } return nil }(), }, }, }, } } // ensureRegistrySecret копирует pull-секрет из namespace оператора в namespace функций. // Вызывается при каждом reconcile — если секрет уже есть, ничего не делает. func (r *FunctionReconciler) ensureRegistrySecret(ctx context.Context, targetNS string) error { // Проверяем что секрет уже есть в целевом namespace existing := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Name: r.RegistrySecret, Namespace: targetNS}, existing); err == nil { return nil // уже есть } else if !errors.IsNotFound(err) { return fmt.Errorf("check secret: %w", err) } // Копируем из namespace оператора src := &corev1.Secret{} if err := r.Get(ctx, client.ObjectKey{Name: r.RegistrySecret, Namespace: r.OperatorNamespace}, src); err != nil { return fmt.Errorf("get source secret from %s: %w", r.OperatorNamespace, err) } copy := &corev1.Secret{ ObjectMeta: metav1.ObjectMeta{ Name: r.RegistrySecret, Namespace: targetNS, }, Type: src.Type, Data: src.Data, } if err := r.Create(ctx, copy); err != nil { if !errors.IsAlreadyExists(err) { return fmt.Errorf("create secret in %s: %w", targetNS, err) } } return nil } // handleDeletion обрабатывает удаление Function: удаляет Deployment, Service, Ingress и убирает finalizer. // ВАЖНО: Namespace sless-fn-{userNS} НЕ удаляется — он принадлежит пользователю на всё время его существования. func (r *FunctionReconciler) handleDeletion(ctx context.Context, fn *slessv1alpha1.Function) (ctrl.Result, error) { deployNS := "sless-fn-" + fn.Namespace dep := &appsv1.Deployment{} if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, dep); err == nil { _ = r.Delete(ctx, dep) } // Если функция удалена в процессе сборки — убиваем kaniko Job. // Без этого Job продолжит работу, займёт CPU/память и запушит образ которым никто не воспользуется. if jobName := fn.Annotations["sless.kube5s.ru/build-job"]; jobName != "" { _ = r.Builder.Cleanup(ctx, jobName) } // Удаляем Service и Ingress — созданы HTTP триггером, но именованы по функции. // Если function_controller не удалит их, Ingress остаётся после destroy → 502. svc := &corev1.Service{} if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, svc); err == nil { _ = r.Delete(ctx, svc) } ing := &netv1.Ingress{} if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, ing); err == nil { _ = r.Delete(ctx, ing) } 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") }