- TriggerReconciler: RequeueAfter 15s когда Function не Ready (ранее зависал без повторного reconcile) - FunctionJobReconciler: RequeueAfter 15s когда Function не Ready - UpdateFunction: добавлена валидация runtime/entrypoint/memory_mb (ранее мог затереть spec нулями при частичном обновлении) - CronJob: curlimages/curl:latest → curlimages/curl:8.5.0 (pin version) - Config: удалён FunctionNamespacePrefix (мёртвое поле, нигде не использовалось) - Invocations endpoint: возвращает 501 вместо пустого списка (SaveInvocation нигде не вызывается — честный ответ клиенту) - Собран образ naeel/sless-operator:v0.1.19
321 lines
12 KiB
Go
321 lines
12 KiB
Go
// Изменено: 2026-03-10
|
||
// TriggerReconciler — контроллер триггеров.
|
||
// HTTP триггер: создаёт Service + Ingress в namespace функции.
|
||
// Cron триггер: создаёт k8s CronJob который периодически вызывает функцию по внутреннему URL.
|
||
|
||
package controllers
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"time"
|
||
|
||
appsv1 "k8s.io/api/apps/v1"
|
||
batchv1 "k8s.io/api/batch/v1"
|
||
corev1 "k8s.io/api/core/v1"
|
||
netv1 "k8s.io/api/networking/v1"
|
||
"k8s.io/apimachinery/pkg/api/errors"
|
||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||
"k8s.io/apimachinery/pkg/runtime"
|
||
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"
|
||
)
|
||
|
||
// TriggerReconciler reconciles a Trigger object
|
||
type TriggerReconciler struct {
|
||
client.Client
|
||
Scheme *runtime.Scheme
|
||
IngressHost string // базовый домен для http триггеров (fallback): fn.kube5s.ru
|
||
// ExternalURL — если задан, URL функции = ExternalURL/fn/{namespace}/{name}
|
||
// (обходит wildcard DNS, трафик идёт через sless-api.kube5s.ru)
|
||
ExternalURL string
|
||
}
|
||
|
||
//+kubebuilder:rbac:groups=sless.kube5s.ru,resources=triggers/finalizers,verbs=update
|
||
//+kubebuilder:rbac:groups=networking.k8s.io,resources=ingresses,verbs=get;list;watch;create;update;patch;delete
|
||
//+kubebuilder:rbac:groups="",resources=services,verbs=get;list;watch;create;update;patch;delete
|
||
//+kubebuilder:rbac:groups=batch,resources=cronjobs,verbs=get;list;watch;create;update;patch;delete
|
||
|
||
const triggerFinalizer = "sless.kube5s.ru/trigger-finalizer"
|
||
|
||
// Reconcile — управляет ресурсами для триггера в зависимости от его типа.
|
||
func (r *TriggerReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
|
||
logger := log.FromContext(ctx)
|
||
|
||
tr := &slessv1alpha1.Trigger{}
|
||
if err := r.Get(ctx, req.NamespacedName, tr); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
return ctrl.Result{}, nil
|
||
}
|
||
return ctrl.Result{}, fmt.Errorf("get trigger: %w", err)
|
||
}
|
||
|
||
// Обработка удаления
|
||
if !tr.DeletionTimestamp.IsZero() {
|
||
return r.handleTriggerDeletion(ctx, tr)
|
||
}
|
||
|
||
// Добавляем finalizer если ещё нет
|
||
if !containsString(tr.Finalizers, triggerFinalizer) {
|
||
tr.Finalizers = append(tr.Finalizers, triggerFinalizer)
|
||
if err := r.Update(ctx, tr); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("add trigger finalizer: %w", err)
|
||
}
|
||
return ctrl.Result{Requeue: true}, nil
|
||
}
|
||
|
||
// Проверяем что Function существует и готова
|
||
fn := &slessv1alpha1.Function{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: tr.Spec.FunctionRef, Namespace: tr.Namespace}, fn); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
tr.Status.Active = false
|
||
tr.Status.Message = "function not found: " + tr.Spec.FunctionRef
|
||
_ = r.Status().Update(ctx, tr)
|
||
return ctrl.Result{}, nil
|
||
}
|
||
return ctrl.Result{}, fmt.Errorf("get function: %w", err)
|
||
}
|
||
if fn.Status.Phase != slessv1alpha1.FunctionPhaseReady {
|
||
tr.Status.Active = false
|
||
tr.Status.Message = "waiting for function to be Ready (current: " + string(fn.Status.Phase) + ")"
|
||
_ = r.Status().Update(ctx, tr)
|
||
// RequeueAfter: опрашиваем каждые 15с пока Function не станет Ready.
|
||
// Watch на Function не добавлен, чтобы не усложнять контроллер на этом этапе.
|
||
return ctrl.Result{RequeueAfter: 15 * time.Second}, nil
|
||
}
|
||
|
||
// Управляем масштабом Deployment функции в зависимости от Enabled.
|
||
// Deployment находится в namespace sless-fn-{userNS}, имя = имя функции.
|
||
// Это позволяет "заморозить" функцию без удаления ресурса.
|
||
deployNS := "sless-fn-" + tr.Namespace
|
||
dep := &appsv1.Deployment{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, dep); err == nil {
|
||
var wantReplicas int32
|
||
if tr.Spec.Enabled {
|
||
wantReplicas = 1
|
||
} else {
|
||
wantReplicas = 0
|
||
}
|
||
// Обновляем только если значение изменилось, чтобы не создавать лишних событий
|
||
if dep.Spec.Replicas == nil || *dep.Spec.Replicas != wantReplicas {
|
||
dep.Spec.Replicas = &wantReplicas
|
||
if err := r.Update(ctx, dep); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("scale deployment replicas: %w", err)
|
||
}
|
||
logger.Info("scaled deployment", "function", fn.Name, "replicas", wantReplicas)
|
||
}
|
||
}
|
||
|
||
// Если триггер выключен — останавливаем reconcile здесь, ресурсы не создаём
|
||
if !tr.Spec.Enabled {
|
||
tr.Status.Active = false
|
||
tr.Status.Message = "disabled (enabled=false)"
|
||
_ = r.Status().Update(ctx, tr)
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
switch tr.Spec.Type {
|
||
case slessv1alpha1.TriggerTypeHTTP:
|
||
logger.Info("reconcile http trigger", "trigger", tr.Name)
|
||
return r.reconcileHTTP(ctx, tr, fn)
|
||
case slessv1alpha1.TriggerTypeCron:
|
||
logger.Info("reconcile cron trigger", "trigger", tr.Name)
|
||
return r.reconcileCron(ctx, tr, fn)
|
||
}
|
||
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// reconcileHTTP создаёт Service для HTTP триггера.
|
||
// Если ExternalURL задан — URL = {ExternalURL}/fn/{namespace}/{name} (прокси через API).
|
||
// Иначе (fallback) — URL = https://{funcName}-{namespace}.{IngressHost} + создаёт Ingress.
|
||
func (r *TriggerReconciler) reconcileHTTP(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) {
|
||
deployNS := "sless-fn-" + tr.Namespace
|
||
|
||
// Service — направляет трафик к Deployment функции (нужен и для прокси, и для Ingress)
|
||
wantSvc := &corev1.Service{
|
||
ObjectMeta: metav1.ObjectMeta{Name: fn.Name, Namespace: deployNS},
|
||
Spec: corev1.ServiceSpec{
|
||
Selector: map[string]string{"app": fn.Name},
|
||
Ports: []corev1.ServicePort{{Port: 8080, Protocol: corev1.ProtocolTCP}},
|
||
},
|
||
}
|
||
existingSvc := &corev1.Service{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingSvc); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
if err := r.Create(ctx, wantSvc); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("create service: %w", err)
|
||
}
|
||
} else {
|
||
return ctrl.Result{}, fmt.Errorf("get service: %w", err)
|
||
}
|
||
}
|
||
|
||
var funcURL string
|
||
if r.ExternalURL != "" {
|
||
// Режим прокси через API сервер — не нужен wildcard DNS
|
||
funcURL = fmt.Sprintf("%s/fn/%s/%s", r.ExternalURL, tr.Namespace, fn.Name)
|
||
} else {
|
||
// Fallback: Ingress с поддоменом (требует wildcard DNS *.IngressHost)
|
||
host := fmt.Sprintf("%s-%s.%s", fn.Name, tr.Namespace, r.IngressHost)
|
||
pathType := netv1.PathTypePrefix
|
||
wantIng := &netv1.Ingress{
|
||
ObjectMeta: metav1.ObjectMeta{
|
||
Name: fn.Name,
|
||
Namespace: deployNS,
|
||
Annotations: map[string]string{
|
||
"kubernetes.io/ingress.class": "nginx",
|
||
},
|
||
},
|
||
Spec: netv1.IngressSpec{
|
||
Rules: []netv1.IngressRule{
|
||
{
|
||
Host: host,
|
||
IngressRuleValue: netv1.IngressRuleValue{
|
||
HTTP: &netv1.HTTPIngressRuleValue{
|
||
Paths: []netv1.HTTPIngressPath{
|
||
{
|
||
Path: "/",
|
||
PathType: &pathType,
|
||
Backend: netv1.IngressBackend{
|
||
Service: &netv1.IngressServiceBackend{
|
||
Name: fn.Name,
|
||
Port: netv1.ServiceBackendPort{Number: 8080},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
}
|
||
existingIng := &netv1.Ingress{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: fn.Name, Namespace: deployNS}, existingIng); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
if err := r.Create(ctx, wantIng); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("create ingress: %w", err)
|
||
}
|
||
} else {
|
||
return ctrl.Result{}, fmt.Errorf("get ingress: %w", err)
|
||
}
|
||
}
|
||
funcURL = "https://" + host
|
||
}
|
||
|
||
tr.Status.Active = true
|
||
tr.Status.URL = funcURL
|
||
tr.Status.Message = ""
|
||
if err := r.Status().Update(ctx, tr); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err)
|
||
}
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// reconcileCron создаёт CronJob который вызывает функцию по HTTP внутри кластера.
|
||
// curl делает POST на внутренний Service функции — это исключает внешний round-trip.
|
||
func (r *TriggerReconciler) reconcileCron(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) {
|
||
deployNS := "sless-fn-" + tr.Namespace
|
||
// Внутренний URL: Service должен быть создан HTTP триггером или заранее
|
||
funcURL := fmt.Sprintf("http://%s.%s.svc.cluster.local:8080", fn.Name, deployNS)
|
||
|
||
wantCJ := &batchv1.CronJob{
|
||
ObjectMeta: metav1.ObjectMeta{
|
||
Name: tr.Name,
|
||
Namespace: tr.Namespace,
|
||
Labels: map[string]string{"managed-by": "sless", "trigger": tr.Name},
|
||
},
|
||
Spec: batchv1.CronJobSpec{
|
||
Schedule: tr.Spec.Schedule,
|
||
JobTemplate: batchv1.JobTemplateSpec{
|
||
Spec: batchv1.JobSpec{
|
||
Template: corev1.PodTemplateSpec{
|
||
Spec: corev1.PodSpec{
|
||
RestartPolicy: corev1.RestartPolicyOnFailure,
|
||
Containers: []corev1.Container{
|
||
{
|
||
// curlimages/curl вызывает функцию по внутреннему адресу
|
||
Name: "invoker",
|
||
Image: "curlimages/curl:latest",
|
||
Command: []string{"curl", "-sf", "-X", "POST", funcURL},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
},
|
||
}
|
||
|
||
existing := &batchv1.CronJob{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, existing); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
if err := r.Create(ctx, wantCJ); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("create cronjob: %w", err)
|
||
}
|
||
} else {
|
||
return ctrl.Result{}, fmt.Errorf("get cronjob: %w", err)
|
||
}
|
||
} else {
|
||
// Синхронизируем расписание если оно изменилось
|
||
existing.Spec.Schedule = tr.Spec.Schedule
|
||
if err := r.Update(ctx, existing); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("update cronjob schedule: %w", err)
|
||
}
|
||
}
|
||
|
||
tr.Status.Active = true
|
||
tr.Status.Message = ""
|
||
if err := r.Status().Update(ctx, tr); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("update trigger status: %w", err)
|
||
}
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// handleTriggerDeletion удаляет ресурсы триггера и убирает finalizer.
|
||
func (r *TriggerReconciler) handleTriggerDeletion(ctx context.Context, tr *slessv1alpha1.Trigger) (ctrl.Result, error) {
|
||
if tr.Spec.Type == slessv1alpha1.TriggerTypeCron {
|
||
cj := &batchv1.CronJob{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: tr.Namespace}, cj); err == nil {
|
||
_ = r.Delete(ctx, cj)
|
||
}
|
||
}
|
||
|
||
// Для HTTP триггеров удаляем Service и Ingress из sless-fn-{namespace}.
|
||
// Имена ресурсов совпадают с именем функции (FunctionRef).
|
||
// Без этого Ingress остаётся висеть после destroy и endpoint возвращает 502.
|
||
if tr.Spec.Type == slessv1alpha1.TriggerTypeHTTP {
|
||
deployNS := "sless-fn-" + tr.Namespace
|
||
fnName := tr.Spec.FunctionRef
|
||
|
||
svc := &corev1.Service{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: fnName, Namespace: deployNS}, svc); err == nil {
|
||
_ = r.Delete(ctx, svc)
|
||
}
|
||
|
||
ing := &netv1.Ingress{}
|
||
if err := r.Get(ctx, client.ObjectKey{Name: fnName, Namespace: deployNS}, ing); err == nil {
|
||
_ = r.Delete(ctx, ing)
|
||
}
|
||
}
|
||
|
||
tr.Finalizers = removeString(tr.Finalizers, triggerFinalizer)
|
||
if err := r.Update(ctx, tr); err != nil {
|
||
return ctrl.Result{}, fmt.Errorf("remove trigger finalizer: %w", err)
|
||
}
|
||
return ctrl.Result{}, nil
|
||
}
|
||
|
||
// SetupWithManager регистрирует контроллер.
|
||
// Watch на Function не добавлен — вместо этого используется RequeueAfter 15s.
|
||
// Это упрощает код на текущем этапе; при росте нагрузки заменить на Watches.
|
||
func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||
return ctrl.NewControllerManagedBy(mgr).
|
||
For(&slessv1alpha1.Trigger{}).
|
||
Complete(r)
|
||
}
|