Files
sless/controllers/function_controller.go

302 lines
13 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Изменено: 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")
}