Files
sless/controllers/function_controller.go
T
“Naeel” 3372cb1983 fix: build logs in status + JSON for all HTTP methods in python runtime
- function_controller.go: KubeClient + getBuildPodLogs → Function.Status.Message
  включает логи pip/kaniko при сбое сборки
- runtimes/python3.11/server.py: PUT/DELETE/PATCH/HEAD обрабатываются как вызовы функции;
  send_error переопределён → JSON вместо HTML 501
- operator.yaml: v0.1.27
- main.go: KubeClient передаётся в FunctionReconciler
2026-03-11 17:03:54 +04:00

437 lines
18 KiB
Go
Raw 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-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")
}