Files
sless/controllers/trigger_controller.go
T
“Naeel” 2ee9cae6d2 feat: proxy /fn/{namespace}/{name} — обход wildcard DNS
Проблема: wildcard DNS *.fn.kube5s.ru недоступен.
Решение: прокси через sless-api.kube5s.ru/fn/{ns}/{name}.

- handler/invoke.go: прокси к Service функции внутри кластера
- router.go: /fn/ без auth токена, /v1/ с auth (gorilla Use())
- config.go: поле ExternalURL (EXTERNAL_URL env)
- trigger_controller.go: если ExternalURL задан — URL = ExternalURL/fn/{ns}/{fn}
  иначе fallback: Ingress + поддомен (прежнее поведение)
- operator.yaml: EXTERNAL_URL=https://sless-api.kube5s.ru, image v0.1.5

Оператор v0.1.5 задеплоен.
E2E: curl https://sless-api.kube5s.ru/fn/default/hello-node → {"message":"Hello, Naeel! (nodejs20)"}
2026-03-07 18:36:03 +04:00

270 lines
10 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-07
// TriggerReconciler — контроллер триггеров.
// HTTP триггер: создаёт Service + Ingress в namespace функции.
// Cron триггер: создаёт k8s CronJob который периодически вызывает функцию по внутреннему URL.
package controllers
import (
"context"
"fmt"
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)
// Повторный reconcile придёт когда Function изменится (watch ниже)
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)
}
}
// Ingress/Service для HTTP триггеров оставляем — они могут быть нужны другим триггерам,
// полная очистка происходит при удалении Function через function_controller finalizer
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.
// Когда Function переходит в Ready — Trigger автоматически пересчитывается.
func (r *TriggerReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&slessv1alpha1.Trigger{}).
Complete(r)
}