Files
sless/controllers/trigger_controller.go

372 lines
15 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-11
// 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)
case slessv1alpha1.TriggerTypeEvent:
logger.Info("reconcile event trigger", "trigger", tr.Name)
return r.reconcileEvent(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.
// CronJob размещается в deployNS (sless-fn-{userNS}), НЕ в user-namespace:
//
// при NetworkPolicy default-deny под в user-ns не может достучаться до Service в sless-fn-ns.
// Размещение CronJob в том же namespace что и Service — гарантирует работу при любой политике.
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: deployNS, // размещаем там же где Service функции
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 вызывает функцию по внутреннему адресу.
// Версия зафиксирована для воспроизводимости — не latest.
Name: "invoker",
Image: "curlimages/curl:8.5.0",
Command: []string{"curl", "-sf", "-X", "POST", funcURL},
},
},
},
},
},
},
},
}
existing := &batchv1.CronJob{}
if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: deployNS}, 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.
// CronJob ищется в deployNS — туда же куда reconcileCron его создаёт.
func (r *TriggerReconciler) handleTriggerDeletion(ctx context.Context, tr *slessv1alpha1.Trigger) (ctrl.Result, error) {
if tr.Spec.Type == slessv1alpha1.TriggerTypeCron {
deployNS := "sless-fn-" + tr.Namespace
cj := &batchv1.CronJob{}
if err := r.Get(ctx, client.ObjectKey{Name: tr.Name, Namespace: deployNS}, 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)
}
// reconcileEvent обрабатывает Trigger{type:event}.
// Оператор не управляет AMQP напрямую — это задача event-dispatcher.
// Здесь: убеждаемся что Service функции существует (dispatcher использует его для POST),
// обновляем статус триггера.
func (r *TriggerReconciler) reconcileEvent(ctx context.Context, tr *slessv1alpha1.Trigger, fn *slessv1alpha1.Function) (ctrl.Result, error) {
if tr.Spec.Queue == "" {
tr.Status.Active = false
tr.Status.Message = "queue is required for type=event"
_ = r.Status().Update(ctx, tr)
return ctrl.Result{}, nil
}
deployNS := "sless-fn-" + tr.Namespace
// Service нужен event-dispatcher для доставки сообщений в функцию по HTTP.
// Имя Service совпадает с именем Function — dispatcher строит URL как
// http://{functionRef}.{deployNS}.svc.cluster.local:8080/
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 for event trigger: %w", err)
}
} else {
return ctrl.Result{}, fmt.Errorf("get service: %w", err)
}
}
tr.Status.Active = true
tr.Status.Message = fmt.Sprintf("listening on queue %q via event-dispatcher", tr.Spec.Queue)
_ = r.Status().Update(ctx, tr)
return ctrl.Result{}, nil
}