302 lines
13 KiB
Go
302 lines
13 KiB
Go
// Изменено: 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")
|
||
}
|